diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/daemon.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 124 |
1 files changed, 109 insertions, 15 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index dc5df81..220a908 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -29,6 +29,7 @@ import base64 import logging import os import signal +import socket import sys import time from dataclasses import asdict @@ -161,6 +162,24 @@ class NodeDaemon(EnrichmentMixin): async def run(self) -> None: log.info("MeshBay Node starting up") + # Set from the first second, not once the node is up: a node waiting for + # its account is exactly the one somebody wants to stop or restart. The + # control API's /api/shutdown is how anything asks -- the one channel + # that reaches a daemon in any session (a service node runs in session 0, + # beyond taskkill and CTRL_BREAK from the user's desktop) and runs the + # real _shutdown() instead of a TerminateProcess. + self._stop_event = asyncio.Event() + loop = asyncio.get_running_loop() + # A beat later, so the HTTP answer to the request leaves first. + self._state["request_shutdown"] = lambda: loop.call_later(0.2, self._stop_event.set) + + # The port before anything else, the token above all: a second instance + # used to write its token over the running node's, then fail to bind + # inside uvicorn's task and exit with the reason on a console nobody sees + # -- leaving the node that did run with a token file matching nothing, so + # every stop, status and restart of it was refused. + ui_sock = bind_control_port(self._config.node.ui_port) + # 1. Keystore keys = load_or_create_keystore( path=self._config.keystore.path, @@ -196,7 +215,7 @@ class NodeDaemon(EnrichmentMixin): log_level="warning", ) ui_server = uvicorn.Server(ui_cfg) - self._tasks.append(asyncio.create_task(ui_server.serve())) + self._tasks.append(asyncio.create_task(ui_server.serve(sockets=[ui_sock]))) log.info("Control API on 127.0.0.1:%d", self._config.node.ui_port) self._tasks.append(asyncio.create_task(self._reap_partial_uploads())) @@ -209,6 +228,10 @@ class NodeDaemon(EnrichmentMixin): async with HubClient(hub_cfg, keys) as hub: self._hub = hub session = await self._login_with_retry(hub) + if session is None: + log.info("Stop requested before the hub sign-in completed") + await self._shutdown() + return self._state["endpoint_hint"] = session.node_id try: @@ -670,13 +693,11 @@ class NodeDaemon(EnrichmentMixin): self._progress_pusher(idx))) # 12. Wait for shutdown - stop_event = asyncio.Event() - loop = asyncio.get_event_loop() + stop_event = self._stop_event console_shutdown_done = None if sys.platform == "win32": - # SIGBREAK: CTRL_BREAK_EVENT, how platform.autostart_end() asks - # a per-user-mode daemon to stop gracefully instead of only - # ever taskkill /F. + # SIGBREAK: Ctrl+Break in the console of a node run by hand. + # Anything else asks through /api/shutdown (request_shutdown). for sig in (signal.SIGINT, signal.SIGTERM, signal.SIGBREAK): signal.signal(sig, lambda *_: stop_event.set()) # CTRL_CLOSE/LOGOFF/SHUTDOWN reach no Python signal at all -- @@ -950,8 +971,29 @@ class NodeDaemon(EnrichmentMixin): gids = list((self._state.get("groups_ctx") or {}).keys()) await self._hub.update_ws_groups(gids) + # Between sign-in attempts while the key is not linked yet. The hub allows + # 10 node sign-ins a minute (api/nodes.py); every 5s was 12, so a node + # waiting for its link rate-limited itself within a minute and spent the + # time it should have been answering the app in 429 back-off. + LINK_WAIT_S = 10 + + async def _pause(self, seconds: float) -> bool: + """Wait, unless a stop is requested first; True when it was.""" + event = getattr(self, "_stop_event", None) + if event is None: + await asyncio.sleep(seconds) + return False + try: + await asyncio.wait_for(event.wait(), timeout=seconds) + return True + except TimeoutError: + return False + async def _login_with_retry(self, hub: HubClient): - """Login to hub, retrying if the node key hasn't been linked yet.""" + """Login to hub, retrying if the node key hasn't been linked yet. + + None when a stop is requested while it waits. + """ import httpx as _httpx while True: try: @@ -971,18 +1013,20 @@ class NodeDaemon(EnrichmentMixin): "Node key not linked. Get it from `meshbay-node " "status` and paste it in Settings > Link Node on %s " "(the desktop client links it automatically). " - "Retrying in 5s...", - self._config.hub.url, + "Retrying in %ds...", + self._config.hub.url, self.LINK_WAIT_S, ) else: self._state["status"] = "waiting_for_account" log.warning( "Hub rejected the node credentials for user %r. " "Register that account on %s first, then link this " - "node's key. Retrying in 5s...", + "node's key. Retrying in %ds...", self._config.hub.username, self._config.hub.url, + self.LINK_WAIT_S, ) - await asyncio.sleep(5) + if await self._pause(self.LINK_WAIT_S): + return None elif e.response.status_code in (429, 500, 502, 503, 504): # Transient: the hub is busy (429 — often this daemon's own # retry storm against the sign-in rate limit), restarting @@ -999,12 +1043,18 @@ class NodeDaemon(EnrichmentMixin): delay = min(max(delay, int(ra)), 300) log.warning("Hub returned %s on login — retrying in %ds", e.response.status_code, delay) - await asyncio.sleep(delay) + if await self._pause(delay): + return None else: raise except Exception as e: + # Unreachable, or an answer that does not verify: the node is up + # and waiting on the hub, and says so rather than keeping a + # status from before. + self._state["status"] = "waiting_for_hub" log.warning("Hub login failed: %s — retrying in 10s", e) - await asyncio.sleep(10) + if await self._pause(10): + return None async def _load_gek( self, @@ -1516,25 +1566,69 @@ class NodeDaemon(EnrichmentMixin): # ── Entry point ─────────────────────────────────────────────────────────────── +class ControlPortTaken(RuntimeError): + pass + + +def bind_control_port(port: int) -> socket.socket: + """The control API's listening socket, or ControlPortTaken. + + Exclusive on Windows, where SO_REUSEADDR would let a second socket share a + port that is in use; SO_REUSEADDR on POSIX only lets a restart reuse a port + left in TIME_WAIT, which is what it is for there. + """ + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + try: + if sys.platform == "win32": + sock.setsockopt(socket.SOL_SOCKET, socket.SO_EXCLUSIVEADDRUSE, 1) + else: + sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + sock.bind(("127.0.0.1", port)) + sock.listen(128) + except OSError as e: + sock.close() + raise ControlPortTaken( + f"the control API's port 127.0.0.1:{port} is taken ({e.strerror or e}) -- " + "a node is probably already running; this one exits. " + "Check with: meshbay-node status") from e + sock.setblocking(False) + return sock + + def main() -> None: args = start() if run(args): return + from meshbay_node.platform import check_media_tools, log_file + if log_file() is not None: + log.info("Logging to %s", log_file()) + + # Printed for whoever ran it, and logged too: a daemon started by the + # service task has no console, and these are why it would refuse to start. cfg = load_config(args.config or DEFAULT_CONFIG_PATH) if not cfg.hub.username: print("Error: hub.username not set in config. Run: meshbay-node init") + log.error("hub.username not set in config. Run: meshbay-node init") sys.exit(1) - from meshbay_node.platform import check_media_tools try: check_media_tools(cfg.node.ffmpeg_path, cfg.node.ffprobe_path) except RuntimeError as e: print(f"Error: {e}") + log.error("%s", e) sys.exit(1) daemon = NodeDaemon(cfg, Path(args.config or DEFAULT_CONFIG_PATH)) - asyncio.run(daemon.run()) + try: + asyncio.run(daemon.run()) + except ControlPortTaken as e: + print(f"Error: {e}") + log.error("%s", e) + sys.exit(2) + except Exception: + log.exception("The node stopped on an unexpected error") + raise if __name__ == "__main__": |