diff options
Diffstat (limited to 'packages')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/cli/dispatch.py | 27 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/cli/lifecycle.py | 154 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 124 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/groups.py | 5 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/platform.py | 188 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ui/app.py | 19 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/conftest.py | 42 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_login_retry_is_resilient.py | 20 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_platform.py | 252 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_windows_service_diagnosis.py | 332 |
10 files changed, 948 insertions, 215 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/cli/dispatch.py b/packages/meshbay-node/src/meshbay_node/cli/dispatch.py index 849d589..fc5712f 100644 --- a/packages/meshbay-node/src/meshbay_node/cli/dispatch.py +++ b/packages/meshbay-node/src/meshbay_node/cli/dispatch.py @@ -28,11 +28,36 @@ def start(): "reset") logging.basicConfig( level=logging.ERROR if quiet else getattr(logging, args.log_level), - format="%(asctime)s %(levelname)-8s %(name)s: %(message)s", + format=LOG_FORMAT, ) + if args.command is None: + add_log_file() return args +LOG_FORMAT = "%(asctime)s %(levelname)-8s %(name)s: %(message)s" + + +def add_log_file() -> None: + """Also log the daemon to platform.log_file(), where there is one.""" + from logging.handlers import RotatingFileHandler + + from meshbay_node.platform import log_file + path = log_file() + if path is None: + return + try: + path.parent.mkdir(parents=True, exist_ok=True) + # delay: nothing is created until there is something to write. + handler = RotatingFileHandler(path, maxBytes=5 * 1024 * 1024, backupCount=3, + encoding="utf-8", delay=True) + except OSError as e: + logging.getLogger(__name__).warning("cannot open the log file %s: %s", path, e) + return + handler.setFormatter(logging.Formatter(LOG_FORMAT)) + logging.getLogger().addHandler(handler) + + # Each verb and what runs it. A command line naming none of them starts the # daemon, which is daemon.main's to do. VERBS = { diff --git a/packages/meshbay-node/src/meshbay_node/cli/lifecycle.py b/packages/meshbay-node/src/meshbay_node/cli/lifecycle.py index 6341561..901c3d5 100644 --- a/packages/meshbay-node/src/meshbay_node/cli/lifecycle.py +++ b/packages/meshbay-node/src/meshbay_node/cli/lifecycle.py @@ -24,28 +24,107 @@ def reload(args) -> None: return +def _await_daemon(cfg, since: float, timeout: float = 30.0) -> dict | None: + """The status of a daemon started at or after `since`, or None. + + The token file is rewritten at every start, so one older than `since` + belongs to the instance that was just stopped — answering from it would + report the old daemon as the new one. + """ + import json + import time + import urllib.request + + token_file = cfg.data_dir / "ui-token" + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + try: + if token_file.stat().st_mtime >= since: + token = token_file.read_text(encoding="utf-8").strip() + url = f"http://127.0.0.1:{cfg.node.ui_port}/api/status?t={token}" + with urllib.request.urlopen(url, timeout=2) as r: + return json.loads(r.read()) + except (OSError, ValueError): + pass + time.sleep(0.5) + return None + + +def _running_status(cfg) -> dict | None: + """The status of whatever node answers right now, or None.""" + return _await_daemon(cfg, since=0, timeout=0.1) + + +def _stop_node(cfg) -> str: + """Stop the node, whichever session it runs in: through its own control API + first, so it shuts down properly, then by force. Returns how it went.""" + import time + + from meshbay_node.platform import ( + autostart_end, + request_graceful_stop, + service_end, + service_state, + service_status, + ) + if request_graceful_stop(cfg.data_dir, cfg.node.ui_port): + how = "stopped" + elif _running_status(cfg) is None: + how = "not running" + else: + how = "stopped (forced)" + if service_status()["installed"]: + # Also leaves the task "Ready": a /run while it still reads "Running" + # is dropped (MultipleInstances IgnoreNew), leaving no node at all. + service_end() + deadline = time.monotonic() + 15 + while service_state().lower() == "running" and time.monotonic() < deadline: + time.sleep(0.25) + autostart_end() # a node in this session that would not stop + return how + + +def _start_and_confirm(args, action: str) -> None: + """Start the node the way this machine is set up to, and report what + actually answered -- never "started" about a node nobody checked.""" + import time + + from meshbay_node import __version__ + from meshbay_node.platform import autostart_run, log_file, service_run, service_status + + cfg = load_config(args.config or DEFAULT_CONFIG_PATH) + if action == "restart": + _stop_node(cfg) + else: + already = _running_status(cfg) + if already is not None: + print(f"already running — node {already.get('version', '?')}, " + f"{already.get('status', '?')}") + return + since = time.time() - 1 + service = service_status()["installed"] + try: + service_run() if service else autostart_run() + except RuntimeError as e: + print(f"Could not start the node: {e}") + sys.exit(1) + status = _await_daemon(cfg, since) + if status is None: + print("no node answered within 30s of being started" + + (" by the service task." if service else ".")) + print(f"See the node's log: {log_file()}") + sys.exit(1) + running = status.get("version", "?") + print(f"{'restarted' if action == 'restart' else 'started'} — node {running}, " + f"{status.get('status', '?')}") + if running != __version__: + print(f"warning: that is node {running} but this is {__version__} — " + "its files were not replaced, or another copy was started.") + + def restart_daemon(args) -> None: if sys.platform == "win32": - from meshbay_node.platform import ( - autostart_end, - autostart_run, - service_end, - service_run, - service_status, - ) - if service_status()["installed"]: - service_end() - service_run() - print("restarted the node (service task)") - return - autostart_end() # kill whatever is running now - try: - autostart_run() - except RuntimeError as e: - print(f"Could not restart: {e}. Stop the daemon (Ctrl+C) and " - "relaunch it from where meshbay-node is on PATH.") - sys.exit(1) - print("restarted the node") + _start_and_confirm(args, "restart") return _systemctl_user( "restart", "meshbay-node", @@ -64,7 +143,14 @@ def autostart(args) -> None: "'systemctl --user enable --now meshbay-node'.") sys.exit(1) sub = args.subcommand or "status" + service_mode = _plat.service_status()["installed"] if sub == "install": + if service_mode: + # Both would start the node: at boot, then again at sign-in. + print("The node already runs as a background service. Remove it " + "first (meshbay-node service remove, elevated) to start it at " + "sign-in instead.") + sys.exit(1) _plat.autostart_install() print("Installed the Startup launcher — meshbay-node starts at " "each sign-in (no window, no admin).") @@ -73,19 +159,15 @@ def autostart(args) -> None: _plat.autostart_remove() print("Removed the Startup launcher.") elif sub == "start": - try: - _plat.autostart_run() - except RuntimeError as e: - print(f"Could not start: {e}") - sys.exit(1) - print("started") + _start_and_confirm(args, "start") elif sub == "stop": - _plat.autostart_end() - print("stopped") + print(_stop_node(load_config(args.config or DEFAULT_CONFIG_PATH))) elif sub == "status": st = _plat.autostart_status() if st["installed"]: print("autostart installed — runs meshbay-node at sign-in") + elif service_mode: + print("autostart not used — the node runs as a background service") else: print("autostart not installed — meshbay-node autostart install") else: @@ -102,6 +184,9 @@ def service(args) -> None: sys.exit(1) sub = args.subcommand or "status" if sub == "install": + # A node already running in this session holds the control API's port: + # the service's own would exit at once, leaving the old one in charge. + _stop_node(load_config(args.config or DEFAULT_CONFIG_PATH)) try: _plat.service_install() except RuntimeError as e: @@ -114,14 +199,19 @@ def service(args) -> None: "signed in yet (no password stored).") print("Start it now with: meshbay-node service start") elif sub == "remove": + # Deleting a task does not end its running instance: stop the node + # first, or it runs on in session 0 with nothing left to stop it. + _stop_node(load_config(args.config or DEFAULT_CONFIG_PATH)) _plat.service_remove() print(f"Removed the {_plat.TASK_NAME!r} scheduled task.") elif sub == "start": - _plat.service_run() - print("started") + if not _plat.service_status()["installed"]: + print("service not installed — meshbay-node service install " + "(needs an elevated prompt)") + sys.exit(1) + _start_and_confirm(args, "start") elif sub == "stop": - _plat.service_end() - print("stopped") + print(_stop_node(load_config(args.config or DEFAULT_CONFIG_PATH))) elif sub == "status": st = _plat.service_status() if st["installed"]: 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__": diff --git a/packages/meshbay-node/src/meshbay_node/ops/groups.py b/packages/meshbay-node/src/meshbay_node/ops/groups.py index 905aa44..1d8003c 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/groups.py +++ b/packages/meshbay-node/src/meshbay_node/ops/groups.py @@ -122,7 +122,10 @@ async def list_groups(state: dict) -> dict: "peers": sum(1 for p in peers.values() if p.get("group_id") == gid), }) roster = state.get("roster") - has_operator = False + # None, not False, until the roster is published (after the hub login): a + # node still waiting for its account has not lost its operator, and saying + # "no operator paired" sent the reader after a pairing that was intact. + has_operator = None if roster: members = await roster.list_members() has_operator = any(m["role"] == "operator" and m["status"] == "active" diff --git a/packages/meshbay-node/src/meshbay_node/platform.py b/packages/meshbay-node/src/meshbay_node/platform.py index 7db251b..09397a2 100644 --- a/packages/meshbay-node/src/meshbay_node/platform.py +++ b/packages/meshbay-node/src/meshbay_node/platform.py @@ -4,7 +4,6 @@ import asyncio import logging import os import shutil -import signal import subprocess import sys import threading @@ -74,6 +73,17 @@ def state_dir() -> Path: return Path.home() / ".local" / "state" / "meshbay" +def log_file() -> Path | None: + """Where the daemon logs on Windows; None elsewhere, where journald has it. + + A daemon started by the service task or the Startup launcher has no console: + its stderr goes nowhere, and a node that fails there fails in silence. + """ + if sys.platform == "win32": + return state_dir() / "node.log" + return None + + # ── Packaged defaults ──────────────────────────────────────────────────────── @@ -211,35 +221,71 @@ def _startup_vbs() -> Path: / "Startup" / "MeshBay Node.vbs") -def _pid_file() -> Path: - """Where autostart_run() records the pid it spawned, for autostart_end() - to signal later -- possibly from a different process (a new Electron - session, or a fresh CLI invocation), so this cannot be an in-memory - handle.""" - return state_dir() / "node.pid" -def _pid_is_meshbay_node(pid: int) -> bool: - """True if `pid` is currently running *and* is meshbay-node.exe. Guards - against a stale pidfile whose pid Windows has since handed to an - unrelated process -- autostart_end() would otherwise send CTRL_BREAK_EVENT - to whatever that is instead.""" - r = subprocess.run(["tasklist", "/FI", f"PID eq {pid}", "/NH"], +def _pid_alive(pid: int) -> bool: + """Whether a process with this pid exists, in any session (tasklist lists + session 0 too, where a service-mode daemon runs; opening it would not).""" + r = subprocess.run(["tasklist", "/FI", f"PID eq {pid}", "/NH", "/FO", "CSV"], capture_output=True, text=True) - return "meshbay-node.exe" in r.stdout.lower() + return f'"{pid}"' in r.stdout def _node_exe() -> str | None: - """Best guess at the meshbay-node launcher: PATH first, then next to the - interpreter (a venv's Scripts/ dir, or a bundled runtime), then argv[0].""" - found = shutil.which("meshbay-node") - if found: - return found + """The meshbay-node to launch: this very one when frozen, then the launcher + next to the interpreter (a venv's Scripts/), then argv[0], then PATH. + + PATH was first, and a machine with two copies on it -- an MSIX and an NSIS + install, an older one left behind -- had `autostart start` from one launch + the other, a different version, found by testing a real install. + """ + if getattr(sys, "frozen", False): + return sys.executable for cand in (Path(sys.executable).parent / "meshbay-node.exe", Path(sys.argv[0])): if cand.name.lower().startswith("meshbay-node") and cand.exists(): return str(cand.resolve()) - return None + return shutil.which("meshbay-node") + + +def request_graceful_stop(data_dir: Path, ui_port: int, timeout: float = 15.0) -> bool: + """Ask the running daemon to stop through its control API, and wait until + its process is gone. True once it is; False when nothing answered or it + did not exit in time -- the caller then forces it. + + The one stop that reaches a daemon in any session with no elevation, and + runs its _shutdown() (WebRTC sessions closed, transcodes stopped). CTRL_BREAK + did neither across sessions, and aimed at a process on another console it + reached every process on the caller's own -- the CLI killed itself. + """ + import json + import urllib.request + + try: + token = (data_dir / "ui-token").read_text(encoding="utf-8").strip() + except OSError: + return False + base = f"http://127.0.0.1:{ui_port}" + try: + with urllib.request.urlopen(f"{base}/api/status?t={token}", timeout=3) as r: + pid = json.loads(r.read()).get("pid") + req = urllib.request.Request(f"{base}/api/shutdown?t={token}", method="POST", data=b"") + with urllib.request.urlopen(req, timeout=3): + pass + except (OSError, ValueError): + return False + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if pid is None: + try: + urllib.request.urlopen(f"{base}/api/status?t={token}", timeout=1) + except OSError: + return True + elif not _pid_alive(int(pid)): + return True + time.sleep(0.3) + log.warning("the node did not exit within %.0fs of being asked to stop", timeout) + return False def autostart_status() -> dict: @@ -283,69 +329,43 @@ def autostart_run() -> None: exe = _node_exe() if not exe: raise RuntimeError("cannot locate the meshbay-node launcher") - # CREATE_NEW_PROCESS_GROUP, not DETACHED_PROCESS: still no visible window - # (CREATE_NO_WINDOW), but the child keeps a console object of its own and - # becomes the root of its own process group -- what autostart_end() needs - # to target it with CTRL_BREAK_EVENT instead of only ever a hard taskkill. - # DETACHED_PROCESS has no console at all, so nothing could be signalled. - proc = subprocess.Popen([exe], creationflags=0x00000200 | 0x08000000, - close_fds=True) - try: - pid_file = _pid_file() - pid_file.parent.mkdir(parents=True, exist_ok=True) - pid_file.write_text(str(proc.pid), encoding="utf-8") - except OSError: - pass # best effort -- autostart_end() falls back to taskkill by image name + # CREATE_NO_WINDOW: a hidden console of its own, which still delivers + # CTRL_LOGOFF/SHUTDOWN to install_console_close_handler, so signing out + # stops the node properly. close_fds: nothing of the caller's is inherited + # -- the desktop app starts its node through here for that reason, since a + # child of Electron inherited Electron's sockets and held them after it quit. + subprocess.Popen([exe], creationflags=0x00000200 | 0x08000000, close_fds=True) -# How long autostart_end() waits for a graceful CTRL_BREAK_EVENT stop before -# giving up and force-killing. A chosen grace period, not an OS-enforced one -# (unlike the ~5 s Windows itself allows a CTRL_CLOSE/LOGOFF/SHUTDOWN handler, -# see install_console_close_handler below -- CTRL_BREAK carries no such ceiling). -_GRACEFUL_STOP_TIMEOUT_SECS = 5.0 +NODE_IMAGE = "meshbay-node.exe" def autostart_end() -> None: """ - Stop the running daemon. + Force-stop a daemon running in this session: the fallback after + request_graceful_stop(), which callers try first. + + It used to send CTRL_BREAK_EVENT to the pid autostart_run() recorded. That + process has a console of its own (CREATE_NO_WINDOW), so it is not a group + on the caller's console -- and GenerateConsoleCtrlEvent aimed at such a pid + reaches every process on the caller's console instead: `autostart stop` + and `restart-daemon` killed themselves, the latter before restarting + anything (found by running both against a real install). + `taskkill /F` is TerminateProcess: cannot be caught, cannot reach session 0. - Tries a graceful stop first: CTRL_BREAK_EVENT to the pid autostart_run() - recorded. Because that process is the root of its own group - (CREATE_NEW_PROCESS_GROUP), daemon.py's own SIGBREAK handler turns this - into the same stop_event.set() SIGINT/SIGTERM already use, running the - real _shutdown() -- closes WebRTC sessions, kills any in-flight ffmpeg - transcode. Falls back to a hard `taskkill /F`, by image name, when there - is no pidfile, the recorded process is already gone, or it does not exit - within the grace period -- same as before this existed, just no longer - the only path. `taskkill /F` itself is TerminateProcess and cannot be made - graceful; nothing can catch it, on any OS. + Never this process or its parent: the CLI is meshbay-node.exe too (and a + venv's launcher is its parent), and `taskkill /IM meshbay-node.exe` killed + the very command that ran it -- `autostart stop` and `restart-daemon` + exited 1 in silence, the latter with the node down. + `/T` because a venv's meshbay-node.exe is a launcher whose python.exe child + is the daemon, and a frozen daemon's children are its ffmpeg transcodes. """ if not autostart_supported(): return - pid_file = _pid_file() - try: - pid = int(pid_file.read_text(encoding="utf-8").strip()) - except (OSError, ValueError): - pid = None - if pid is not None and not _pid_is_meshbay_node(pid): - pid = None # stale pidfile -- Windows may have reused the pid since - if pid is not None: - try: - os.kill(pid, signal.CTRL_BREAK_EVENT) - except OSError: - pid = None # already gone, or never existed - else: - deadline = time.monotonic() + _GRACEFUL_STOP_TIMEOUT_SECS - while time.monotonic() < deadline: - if not _pid_is_meshbay_node(pid): - pid_file.unlink(missing_ok=True) - return - time.sleep(0.2) - log.warning("pid %d did not exit within %.1fs of CTRL_BREAK_EVENT, " - "falling back to taskkill /F", pid, _GRACEFUL_STOP_TIMEOUT_SECS) - pid_file.unlink(missing_ok=True) - subprocess.run(["taskkill", "/IM", "meshbay-node.exe", "/F"], - capture_output=True) + argv = ["taskkill", "/F", "/T", "/IM", NODE_IMAGE] + for pid in {os.getpid(), os.getppid()}: + argv += ["/FI", f"PID ne {pid}"] + subprocess.run(argv, capture_output=True) # ── Service mode (Windows, opt-in at install time) ─────────────────────────── @@ -393,6 +413,15 @@ def autostart_end() -> None: TASK_NAME = "MeshBay Node" # the Scheduled Task's own name +# Task Scheduler's defaults end a task after 72 hours, never start it on +# battery and stop it when the power cable comes out — a node on a laptop was +# down for any of the three. Kept identical to packaging/win/service.ps1. +SERVICE_TASK_SETTINGS = ( + "New-ScheduledTaskSettingsSet -ExecutionTimeLimit ([TimeSpan]::Zero) " + "-AllowStartIfOnBatteries -DontStopIfGoingOnBatteries " + "-MultipleInstances IgnoreNew -StartWhenAvailable" +) + def service_supported() -> bool: return sys.platform == "win32" @@ -424,6 +453,11 @@ def service_status() -> dict: return {"installed": True, "state": state} +def service_state() -> str: + """Task Scheduler's word for the task ('Running', 'Ready', ...), '' if absent.""" + return service_status()["state"] + + def service_install(exe: str | None = None) -> None: """ Register the boot-time Scheduled Task with an S4U logon. Needs admin — @@ -458,8 +492,9 @@ def service_install(exe: str | None = None) -> None: "$t = New-ScheduledTaskTrigger -AtStartup; " "$p = New-ScheduledTaskPrincipal -UserId $env:MESHBAY_SVC_USER " "-LogonType S4U -RunLevel Limited; " + f"$s = {SERVICE_TASK_SETTINGS}; " "Register-ScheduledTask -TaskName $env:MESHBAY_SVC_TASK " - "-Action $a -Trigger $t -Principal $p -Force | Out-Null" + "-Action $a -Trigger $t -Principal $p -Settings $s -Force | Out-Null" ) r = subprocess.run( ["powershell", "-NoProfile", "-NonInteractive", "-Command", ps_script], @@ -473,6 +508,7 @@ def service_remove() -> None: """Delete the Scheduled Task if present. Needs admin; silent otherwise (mirrors autostart_remove — nothing to report if it was never installed).""" if service_supported(): + _schtasks("/end", "/tn", TASK_NAME) # deleting does not end the instance _schtasks("/delete", "/tn", TASK_NAME, "/f") @@ -498,8 +534,8 @@ def service_end() -> None: # _shutdown(), no closed WebRTC sessions, no killed ffmpeg. Covers: closing # the console window of an interactively-run `meshbay-node run`, user logoff, # system shutdown. Does NOT cover `taskkill /F` -- TerminateProcess is -# uncatchable on any OS, the same as SIGKILL; see autostart_end() for how the -# Node page's Stop button avoids relying on it instead. +# uncatchable on any OS, the same as SIGKILL; request_graceful_stop() is how +# every Stop avoids relying on it. _CONSOLE_HANDLER_REFS: list = [] # ctypes callbacks must be kept referenced or they may be freed diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index 090eb21..50beb26 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -13,6 +13,7 @@ no server-rendered UI: the Node page ships in the desktop client (see """ import logging +import os from fastapi import FastAPI, HTTPException, Query from fastapi.responses import JSONResponse @@ -146,6 +147,9 @@ def create_ui_app(state: dict) -> FastAPI: return { "version": __version__, + # Whoever asks it to stop waits for this process to be gone, not only + # for the API to close: _shutdown() runs on after the API has closed. + "pid": os.getpid(), "status": status, "needs": needs, "hub_url": state.get("hub_url", ""), @@ -542,6 +546,21 @@ def create_ui_app(state: dict) -> FastAPI: # seconds, on a real library) — see ops.start_reload for why. return await _op(lambda: ops.start_reload(state)) + # ── Stop ───────────────────────────────────────────────────────────── + + @app.post("/api/shutdown") + async def shutdown(): + # The graceful stop every front door tries first (the desktop app, the + # CLI, the installer): it reaches a node in any session -- a service + # node runs in session 0, where taskkill and CTRL_BREAK from the user's + # desktop do not -- and runs the node's own _shutdown(). Behind the + # per-run token like everything here. + request_shutdown = state.get("request_shutdown") + if request_shutdown is None: + return JSONResponse({"error": "not ready to stop yet"}, status_code=503) + request_shutdown() + return {"stopping": True} + # ── Node settings (operator only, localhost) ─────────────────────────── @app.get("/api/node-settings") diff --git a/packages/meshbay-node/tests/conftest.py b/packages/meshbay-node/tests/conftest.py index 9ef23a3..32c9b3d 100644 --- a/packages/meshbay-node/tests/conftest.py +++ b/packages/meshbay-node/tests/conftest.py @@ -18,6 +18,48 @@ needs_subprocess = pytest.mark.skipif( ) @pytest.fixture(autouse=True) +def _no_log_file_in_the_developers_profile(monkeypatch): + """A test that runs the daemon's entry point must not open its log file. + + On Windows that file is %LOCALAPPDATA%\\meshbay\\state\\node.log -- the + developer's own node's log. The first run of the suite with it appended + every later test's log lines there, through a handler left on the root + logger. A test about the log file points log_file() somewhere of its own. + """ + import logging + + from meshbay_node import platform as _plat + monkeypatch.setattr(_plat, "log_file", lambda: None) + root = logging.getLogger() + before = list(root.handlers) + yield + for h in root.handlers[:]: + if h not in before and getattr(h, "baseFilename", None): + root.removeHandler(h) + h.close() + + +@pytest.fixture(scope="session") +def _throwaway_profile(tmp_path_factory): + return tmp_path_factory.mktemp("profile") + + +@pytest.fixture(autouse=True) +def _no_test_touches_the_developers_profile(monkeypatch, _throwaway_profile): + """Every per-user location, pointed somewhere nobody lives. + + Tests redirected HOME and nothing else, which isolates nothing on Windows: + Path.home() reads USERPROFILE there and the node's own directories come from + LOCALAPPDATA. `member invite` and `operator pair`, walked by the CLI tests, + wrote their codes into the developer's real node directory and a + ~/.local/share/meshbay nobody had, on every run. A test that sets its own + HOME still can; it just no longer falls through to the real one. + """ + for var in ("HOME", "USERPROFILE", "LOCALAPPDATA", "APPDATA"): + monkeypatch.setenv(var, str(_throwaway_profile / var.lower())) + + +@pytest.fixture(autouse=True) def _restore_media_tool_paths(): """Put `platform`'s resolved ffmpeg/ffprobe paths back after every test. diff --git a/packages/meshbay-node/tests/test_login_retry_is_resilient.py b/packages/meshbay-node/tests/test_login_retry_is_resilient.py index e37a413..39358bb 100644 --- a/packages/meshbay-node/tests/test_login_retry_is_resilient.py +++ b/packages/meshbay-node/tests/test_login_retry_is_resilient.py @@ -14,8 +14,7 @@ import asyncio import httpx import pytest -from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey -from meshbay_node.config import Config, GroupConfig, HubConfig, KeystoreConfig, NodeConfig +from meshbay_node.config import Config, HubConfig, KeystoreConfig, NodeConfig from meshbay_node.daemon import NodeDaemon @@ -96,6 +95,23 @@ async def test_a_401_still_retries_and_stays_alive(tmp_path, monkeypatch): assert daemon._state.get("status") in ("waiting_for_node_key", "waiting_for_account") +def test_a_node_waiting_for_its_link_stays_under_the_hubs_sign_in_limit(): + """At 5s a node waiting for its key made 12 attempts a minute against a + limit of 10: it put itself in 429 back-off within a minute, every minute, + and the desktop app did not recognise that state as a node at all.""" + import re + from pathlib import Path + + from meshbay_node.daemon import NodeDaemon + + nodes_py = (Path(__file__).resolve().parents[2] / "meshbay-hub" / "src" / "meshbay_hub" + / "api" / "nodes.py").read_text(encoding="utf-8") + m = re.search(r'@router\.post\("/auth"\)\s*@limiter\.limit\("(\d+)/minute"\)', nodes_py) + assert m, "the hub's node sign-in limit moved; update this test" + per_minute = int(m.group(1)) + assert 60 / NodeDaemon.LINK_WAIT_S < per_minute + + @pytest.mark.asyncio async def test_a_genuine_client_error_still_raises(tmp_path, monkeypatch): """A 400/422 is a bug, not a transient state — it must not be swallowed.""" diff --git a/packages/meshbay-node/tests/test_platform.py b/packages/meshbay-node/tests/test_platform.py index bebe8d9..8856962 100644 --- a/packages/meshbay-node/tests/test_platform.py +++ b/packages/meshbay-node/tests/test_platform.py @@ -8,6 +8,7 @@ suite happens to run on. import asyncio import os import sys +import time from pathlib import Path from unittest.mock import Mock @@ -168,10 +169,10 @@ def test_autostart_install_refuses_off_windows(monkeypatch): plat.autostart_install(exe="/usr/bin/meshbay-node") -def test_autostart_run_launches_the_resolved_exe_windowless_and_records_its_pid( - win_startup, monkeypatch, tmp_path): +def test_autostart_run_launches_windowless_and_inherits_nothing(win_startup, monkeypatch): + """close_fds: the desktop app starts its node through here because a child + of Electron inherited Electron's sockets and held them after the app quit.""" monkeypatch.setattr(plat, "_node_exe", lambda: r"C:\x\meshbay-node.exe") - monkeypatch.setenv("LOCALAPPDATA", str(tmp_path)) # state_dir() -> pidfile location calls = {} def fake_popen(argv, **kw): @@ -181,13 +182,23 @@ def test_autostart_run_launches_the_resolved_exe_windowless_and_records_its_pid( monkeypatch.setattr(plat.subprocess, "Popen", fake_popen) plat.autostart_run() assert calls["argv"] == [r"C:\x\meshbay-node.exe"] - flags = calls["kw"]["creationflags"] - assert flags & 0x08000000 # CREATE_NO_WINDOW - assert flags & 0x00000200 # CREATE_NEW_PROCESS_GROUP - assert not flags & 0x00000008 # not DETACHED_PROCESS -- that has no - # console at all, so CTRL_BREAK_EVENT - # would have nothing to signal - assert plat._pid_file().read_text(encoding="utf-8") == "4242" + assert calls["kw"]["creationflags"] & 0x08000000 # CREATE_NO_WINDOW + assert not calls["kw"]["creationflags"] & 0x00000008 # a console of its own: + # CTRL_LOGOFF/SHUTDOWN reach install_console_close_handler + assert calls["kw"]["close_fds"] is True + + +def test_the_launcher_is_this_very_executable_when_frozen(monkeypatch, tmp_path): + """PATH came first, and with two copies on it `autostart start` from one + install launched the other -- found on a machine with an MSIX build on PATH.""" + stray = tmp_path / "stray" + stray.mkdir() + (stray / "meshbay-node.exe").write_text("", encoding="utf-8") + monkeypatch.setenv("PATH", str(stray)) + monkeypatch.setattr(plat.sys, "frozen", True, raising=False) + frozen = r"C:\Programs\MeshBay\node-runtime\meshbay-node.exe" + monkeypatch.setattr(plat.sys, "executable", frozen) + assert plat._node_exe() == frozen def test_autostart_run_refuses_off_windows(monkeypatch): @@ -196,100 +207,165 @@ def test_autostart_run_refuses_off_windows(monkeypatch): plat.autostart_run() -# ── Graceful stop (CTRL_BREAK_EVENT + taskkill fallback) ──────────────────── +# ── Stopping ──────────────────────────────────────────────────────────────── # -# autostart_end() references signal.CTRL_BREAK_EVENT, which genuinely does not -# exist in the `signal` module off Windows -- monkeypatching sys.platform -# cannot manufacture it, unlike the pure-Python behaviour tested above. Skip -# rather than mock around it, matching test_configure_event_loop_selector_opt_in. +# The CTRL_BREAK_EVENT path these tests used to pin was the bug: aimed at a +# node with a console of its own, GenerateConsoleCtrlEvent reached every process +# on the caller's console, so `autostart stop` and `restart-daemon` killed +# themselves. The tests mocked os.kill, so they agreed with it by construction. -@pytest.mark.skipif(sys.platform != "win32", - reason="signal.CTRL_BREAK_EVENT exists only on win32") -def test_autostart_end_stops_gracefully_when_ctrl_break_is_enough(monkeypatch, tmp_path): +def test_autostart_end_force_stops_by_image_name(monkeypatch): monkeypatch.setattr(sys, "platform", "win32") - monkeypatch.setenv("LOCALAPPDATA", str(tmp_path)) - plat._pid_file().parent.mkdir(parents=True, exist_ok=True) - plat._pid_file().write_text("4242", encoding="utf-8") - - kill_calls = [] - monkeypatch.setattr(plat.os, "kill", lambda pid, sig: kill_calls.append((pid, sig))) - # Alive (our exe) on the pre-signal check, gone by the first poll after -- - # a plain constant can't tell those two calls apart. - seen = {"n": 0} - - def fake_check(pid): - seen["n"] += 1 - return seen["n"] == 1 - - monkeypatch.setattr(plat, "_pid_is_meshbay_node", fake_check) run_calls = [] monkeypatch.setattr(plat.subprocess, "run", lambda argv, **kw: run_calls.append(argv)) - plat.autostart_end() - - assert kill_calls == [(4242, plat.signal.CTRL_BREAK_EVENT)] - assert run_calls == [] # no taskkill needed - assert not plat._pid_file().exists() + [argv] = run_calls + assert argv[:5] == ["taskkill", "/F", "/T", "/IM", "meshbay-node.exe"] -@pytest.mark.skipif(sys.platform != "win32", - reason="signal.CTRL_BREAK_EVENT exists only on win32") -def test_autostart_end_falls_back_to_taskkill_when_the_pid_never_exits( - monkeypatch, tmp_path): +def test_the_forced_stop_spares_the_command_that_runs_it(monkeypatch): + """The CLI is meshbay-node.exe too: `taskkill /IM meshbay-node.exe` killed + `autostart stop` and `restart-daemon` themselves -- found on a real install.""" monkeypatch.setattr(sys, "platform", "win32") - monkeypatch.setenv("LOCALAPPDATA", str(tmp_path)) - plat._pid_file().parent.mkdir(parents=True, exist_ok=True) - plat._pid_file().write_text("4242", encoding="utf-8") - - monkeypatch.setattr(plat.os, "kill", lambda pid, sig: None) - monkeypatch.setattr(plat, "_pid_is_meshbay_node", lambda pid: True) # never exits - monkeypatch.setattr(plat.time, "sleep", lambda s: None) # don't really wait - clock = iter([0.0, 1.0, 6.0]) # deadline = 0.0 + 5.0; third read is past it - monkeypatch.setattr(plat.time, "monotonic", lambda: next(clock)) run_calls = [] monkeypatch.setattr(plat.subprocess, "run", lambda argv, **kw: run_calls.append(argv)) - plat.autostart_end() + [argv] = run_calls + assert f"PID ne {os.getpid()}" in argv and f"PID ne {os.getppid()}" in argv - assert run_calls == [["taskkill", "/IM", "meshbay-node.exe", "/F"]] - assert not plat._pid_file().exists() +@pytest.mark.skipif(sys.platform != "win32", reason="taskkill") +def test_the_forced_stop_really_spares_its_caller(tmp_path): + """For real: a process with the node's image name runs the forced stop, and + must survive it while another process of that name does not. Under a name of + its own, or it would also kill the developer's running node.""" + import shutil + import subprocess + import textwrap + import uuid -def test_autostart_end_falls_back_to_taskkill_without_a_pidfile(monkeypatch, tmp_path): - """No CTRL_BREAK_EVENT dependency here -- there is no pid to signal, so - this one runs everywhere, same as the pre-existing behaviour it replaces.""" - monkeypatch.setattr(sys, "platform", "win32") - monkeypatch.setenv("LOCALAPPDATA", str(tmp_path)) - run_calls = [] - monkeypatch.setattr(plat.subprocess, "run", - lambda argv, **kw: run_calls.append(argv)) - plat.autostart_end() - assert run_calls == [["taskkill", "/IM", "meshbay-node.exe", "/F"]] + image = f"mb-test-{uuid.uuid4().hex[:8]}.exe" + base = Path(getattr(sys, "_base_executable", sys.executable)) + env = {**os.environ, "PYTHONHOME": sys.base_prefix} + def renamed_python(where: str) -> Path: + d = tmp_path / where + d.mkdir() + shutil.copy(base, d / image) + for dll in base.parent.glob("*.dll"): + shutil.copy(dll, d) + return d / image -def test_autostart_end_ignores_a_stale_pid_reused_by_another_process(monkeypatch, tmp_path): - """The recorded pid is alive but is not meshbay-node.exe -- Windows reused - it after the daemon exited. Must not send CTRL_BREAK_EVENT to whatever - that is; falls straight to taskkill (by image name, so harmless here).""" - monkeypatch.setattr(sys, "platform", "win32") - monkeypatch.setenv("LOCALAPPDATA", str(tmp_path)) - plat._pid_file().parent.mkdir(parents=True, exist_ok=True) - plat._pid_file().write_text("4242", encoding="utf-8") + caller = renamed_python("caller") + victim = renamed_python("victim") # a "daemon" with a "transcode" + child_pid_file = tmp_path / "child.pid" + src = Path(plat.__file__).resolve().parents[1] + script = textwrap.dedent(f""" + import sys; sys.path.insert(0, {str(src)!r}) + from meshbay_node import platform as p + p.NODE_IMAGE = {image!r} + p.autostart_end() + print("survived") + """) + other = subprocess.Popen( + [str(victim), "-c", + "import subprocess, sys; p = subprocess.Popen(['ping', '-n', '300', '127.0.0.1']," + " stdout=subprocess.DEVNULL); open(sys.argv[1], 'w').write(str(p.pid)); p.wait()", + str(child_pid_file)], env=env) + child_pid = 0 + try: + for _ in range(100): + if child_pid_file.exists() and child_pid_file.read_text(encoding="utf-8"): + break + time.sleep(0.1) + child_pid = int(child_pid_file.read_text(encoding="utf-8")) + r = subprocess.run([str(caller), "-c", script], capture_output=True, text=True, + timeout=60, env=env) + assert "survived" in r.stdout, (r.returncode, r.stdout, r.stderr) + assert other.wait(timeout=10) is not None, "the other process must be stopped" + assert not plat._pid_alive(child_pid), "and its children with it" + finally: + if other.poll() is None: + other.kill() + if child_pid: + subprocess.run(["taskkill", "/F", "/PID", str(child_pid)], capture_output=True) - monkeypatch.setattr(plat, "_pid_is_meshbay_node", lambda pid: False) - kill_calls = [] - monkeypatch.setattr(plat.os, "kill", lambda pid, sig: kill_calls.append((pid, sig))) - run_calls = [] - monkeypatch.setattr(plat.subprocess, "run", - lambda argv, **kw: run_calls.append(argv)) - plat.autostart_end() +def test_nothing_signals_a_console_group_any_more(): + """Read as code, not text: the docstrings say why CTRL_BREAK went.""" + import ast + import inspect + tree = ast.parse(inspect.getsource(plat)) + for node in ast.walk(tree): + assert not (isinstance(node, ast.Attribute) and node.attr == "CTRL_BREAK_EVENT") + if isinstance(node, ast.Call) and isinstance(node.func, ast.Attribute): + assert not (node.func.attr == "kill" and getattr(node.func.value, "id", "") == "os") + + +def _serve_node(tmp_path, pid): + """A control API that answers status and shutdown, as the daemon's does.""" + import json + import threading + from http.server import BaseHTTPRequestHandler, HTTPServer + + asked = [] + + class H(BaseHTTPRequestHandler): + def _send(self, body): + data = json.dumps(body).encode() + self.send_response(200) + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def do_GET(self): # noqa: N802 + self._send({"version": "x", "status": "running", "pid": pid}) + + def do_POST(self): # noqa: N802 + asked.append(self.path) + self._send({"stopping": True}) + + def log_message(self, *a): + pass + + server = HTTPServer(("127.0.0.1", 0), H) + threading.Thread(target=server.serve_forever, daemon=True).start() + (tmp_path / "ui-token").write_text("tok", encoding="utf-8") + return server, asked + + +def test_a_graceful_stop_asks_the_node_and_waits_for_its_process(monkeypatch, tmp_path): + server, asked = _serve_node(tmp_path, pid=4242) + alive = iter([True, True, False]) + monkeypatch.setattr(plat, "_pid_alive", lambda pid: next(alive)) + try: + assert plat.request_graceful_stop(tmp_path, server.server_address[1], timeout=10) + finally: + server.shutdown() + assert asked == ["/api/shutdown?t=tok"] + + +def test_a_node_that_does_not_exit_is_reported_for_the_caller_to_force(monkeypatch, tmp_path): + server, _ = _serve_node(tmp_path, pid=4242) + monkeypatch.setattr(plat, "_pid_alive", lambda pid: True) + try: + assert not plat.request_graceful_stop(tmp_path, server.server_address[1], timeout=1) + finally: + server.shutdown() + + +def test_nothing_to_stop_when_no_node_answers(tmp_path): + (tmp_path / "ui-token").write_text("tok", encoding="utf-8") + assert not plat.request_graceful_stop(tmp_path, 1, timeout=1) + - assert kill_calls == [] # never signalled the reused pid - assert run_calls == [["taskkill", "/IM", "meshbay-node.exe", "/F"]] - assert not plat._pid_file().exists() +@pytest.mark.skipif(sys.platform != "win32", reason="tasklist") +def test_pid_alive_sees_this_process_and_not_a_dead_one(): + import os + assert plat._pid_alive(os.getpid()) + assert not plat._pid_alive(9_999_991) def test_autostart_end_is_a_noop_off_windows(monkeypatch): @@ -413,7 +489,7 @@ def test_frozen_build_finds_default_env_beside_the_executable(monkeypatch, tmp_p """Where build-node-runtime.ps1 puts it, alongside ffmpeg.""" exe = tmp_path / "meshbay-node.exe" exe.write_bytes(b"") - (tmp_path / "default.env").write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJtest\n") + (tmp_path / "default.env").write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJtest\n", encoding="utf-8") monkeypatch.setattr(sys, "frozen", True, raising=False) monkeypatch.setattr(sys, "executable", str(exe)) assert plat.packaged_default_env() == tmp_path / "default.env" @@ -434,7 +510,7 @@ def test_load_node_env_sets_names(monkeypatch, tmp_path): "\n" "MESHBAY_TMDB_DEFAULT_TOKEN=eyJloaded\n" 'QUOTED="value"\n' - ) + , encoding="utf-8") monkeypatch.delenv("MESHBAY_TMDB_DEFAULT_TOKEN", raising=False) monkeypatch.delenv("QUOTED", raising=False) assert plat.load_node_env(tmp_path) == 2 @@ -445,7 +521,7 @@ def test_load_node_env_sets_names(monkeypatch, tmp_path): def test_load_node_env_does_not_override_the_environment(monkeypatch, tmp_path): """systemd may have loaded the same file already, and an operator export must win over a packaged default.""" - (tmp_path / "node.env").write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJfromfile\n") + (tmp_path / "node.env").write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJfromfile\n", encoding="utf-8") monkeypatch.setenv("MESHBAY_TMDB_DEFAULT_TOKEN", "eyJfromenv") assert plat.load_node_env(tmp_path) == 0 assert os.environ["MESHBAY_TMDB_DEFAULT_TOKEN"] == "eyJfromenv" @@ -460,7 +536,7 @@ def test_load_node_env_reads_the_packaged_default_without_a_node_env(monkeypatch """A node onboarded by the desktop client has no node.env: nothing ran `init` to copy one. The packaged token must reach it anyway.""" src = tmp_path / "default.env" - src.write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJpackaged\n") + src.write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJpackaged\n", encoding="utf-8") monkeypatch.setattr(plat, "packaged_default_env", lambda: src) monkeypatch.delenv("MESHBAY_TMDB_DEFAULT_TOKEN", raising=False) cfg = tmp_path / "config" @@ -472,12 +548,12 @@ def test_load_node_env_reads_the_packaged_default_without_a_node_env(monkeypatch def test_load_node_env_prefers_the_operator_node_env(monkeypatch, tmp_path): src = tmp_path / "default.env" - src.write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJpackaged\n") + src.write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJpackaged\n", encoding="utf-8") monkeypatch.setattr(plat, "packaged_default_env", lambda: src) monkeypatch.delenv("MESHBAY_TMDB_DEFAULT_TOKEN", raising=False) cfg = tmp_path / "config" cfg.mkdir() - (cfg / "node.env").write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJoperator\n") + (cfg / "node.env").write_text("MESHBAY_TMDB_DEFAULT_TOKEN=eyJoperator\n", encoding="utf-8") plat.load_node_env(cfg) assert os.environ["MESHBAY_TMDB_DEFAULT_TOKEN"] == "eyJoperator" diff --git a/packages/meshbay-node/tests/test_windows_service_diagnosis.py b/packages/meshbay-node/tests/test_windows_service_diagnosis.py new file mode 100644 index 0000000..84ba0b6 --- /dev/null +++ b/packages/meshbay-node/tests/test_windows_service_diagnosis.py @@ -0,0 +1,332 @@ +""" +What a Windows node says about itself when something is wrong. + +A 0.16 upgrade in service mode left the previous version's node running, and +every message on the way out pointed elsewhere: `meshbay-node service start` +printed "started" for a node that had not, the app said "the node started but +could not link to your hub account" about one that had stopped answering, the +Node page said "No operator paired" about one whose pairing was intact, and the +daemon itself had logged nowhere at all -- Task Scheduler discards its stderr. +""" + +import argparse +import json +import logging +import os +import threading +import time +from http.server import BaseHTTPRequestHandler, HTTPServer + +import pytest +from meshbay_node import __version__ +from meshbay_node import platform as plat +from meshbay_node.cli import dispatch, lifecycle + +# ── the daemon's log file ────────────────────────────────────────────────── + + +@pytest.fixture +def root_logger(): + """add_log_file changes the root logger: pin it and put it back.""" + root = logging.getLogger() + handlers, level, disabled = list(root.handlers), root.level, root.disabled + root.setLevel(logging.INFO) + root.disabled = False + yield root + for h in root.handlers: + if h not in handlers: + h.close() + root.handlers[:] = handlers + root.setLevel(level) + root.disabled = disabled + + +def test_the_daemon_logs_to_a_file_where_there_is_no_console(tmp_path, monkeypatch, root_logger): + log_path = tmp_path / "state" / "node.log" + monkeypatch.setattr(plat, "log_file", lambda: log_path) + + dispatch.add_log_file() + logging.getLogger("meshbay_node.daemon").warning("hub login failed: probe") + for h in root_logger.handlers: + h.flush() + + assert log_path.exists(), "the log file's directory must be created, not assumed" + assert "hub login failed: probe" in log_path.read_text(encoding="utf-8") + + +def test_no_log_file_where_journald_has_the_output(monkeypatch, root_logger): + monkeypatch.setattr(plat, "log_file", lambda: None) + before = list(root_logger.handlers) + dispatch.add_log_file() + assert root_logger.handlers == before + + +def test_the_log_file_is_under_localappdata_on_windows(monkeypatch, tmp_path): + monkeypatch.undo() # the conftest keeps log_file() away from the real profile + monkeypatch.setattr(plat.sys, "platform", "win32") + monkeypatch.setenv("LOCALAPPDATA", str(tmp_path)) + assert plat.log_file() == tmp_path / "meshbay" / "state" / "node.log" + monkeypatch.setattr(plat.sys, "platform", "linux") + assert plat.log_file() is None + + +def test_only_the_daemon_run_opens_the_log_file(monkeypatch): + """CLI verbs print reports; appending them to the daemon's log would bury + the daemon's own lines.""" + calls = [] + monkeypatch.setattr(dispatch, "add_log_file", lambda: calls.append(1)) + monkeypatch.setattr(plat, "load_node_env", lambda _d: 0) + for argv, expected in ((["meshbay-node", "status"], 0), (["meshbay-node"], 1)): + calls.clear() + monkeypatch.setattr("sys.argv", argv) + dispatch.start() + assert len(calls) == expected, argv + + +# ── `service start` says what happened ────────────────────────────────────── + + +class _Status(BaseHTTPRequestHandler): + body = {} + + def do_GET(self): # noqa: N802 + data = json.dumps(self.body).encode() + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def log_message(self, *_a): + pass + + +@pytest.fixture +def fake_daemon(tmp_path): + server = HTTPServer(("127.0.0.1", 0), _Status) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + data_dir = tmp_path / "data" + data_dir.mkdir() + cfg = argparse.Namespace( + data_dir=data_dir, + node=argparse.Namespace(ui_port=server.server_address[1])) + yield cfg + server.shutdown() + + +def test_a_token_from_the_stopped_instance_is_not_taken_for_the_new_one(fake_daemon): + _Status.body = {"version": "0.16.0", "status": "running"} + token = fake_daemon.data_dir / "ui-token" + token.write_text("t", encoding="utf-8") + old = time.time() - 60 + os.utime(token, (old, old)) + + assert lifecycle._await_daemon(fake_daemon, since=time.time() - 1, timeout=1.5) is None + + +def test_a_daemon_that_started_is_seen(fake_daemon): + _Status.body = {"version": "0.16.0", "status": "waiting_for_account"} + since = time.time() - 1 + (fake_daemon.data_dir / "ui-token").write_text("t", encoding="utf-8") + + got = lifecycle._await_daemon(fake_daemon, since=since, timeout=5) + assert got == _Status.body + + +def _service_start(monkeypatch, awaited, capsys): + monkeypatch.setattr(plat, "service_status", lambda: {"installed": True, "state": "Ready"}) + monkeypatch.setattr(plat, "service_run", lambda: None) + monkeypatch.setattr(plat, "service_end", lambda: None) + monkeypatch.setattr(plat, "log_file", lambda: "C:/x/meshbay/state/node.log") + monkeypatch.setattr(lifecycle, "load_config", lambda _p: object()) + monkeypatch.setattr(lifecycle, "_running_status", lambda _cfg: None) + monkeypatch.setattr(lifecycle, "_await_daemon", lambda _cfg, _since, **_kw: awaited) + monkeypatch.setattr(lifecycle.sys, "platform", "win32") + lifecycle.service(argparse.Namespace(subcommand="start", config=None)) + return capsys.readouterr().out + + +def test_service_start_does_not_report_a_node_that_never_answered(monkeypatch, capsys): + with pytest.raises(SystemExit) as exc: + _service_start(monkeypatch, None, capsys) + assert exc.value.code == 1 + out = capsys.readouterr().out + assert not out.startswith("started"), "it must not claim the node started" + assert "no node answered" in out + assert "node.log" in out, "and must say where to look" + + +def test_service_start_names_the_version_that_answered(monkeypatch, capsys): + out = _service_start(monkeypatch, {"version": __version__, "status": "running"}, capsys) + assert f"node {__version__}" in out and "warning" not in out + + +def test_service_start_warns_when_the_service_runs_an_older_node(monkeypatch, capsys): + out = _service_start(monkeypatch, {"version": "0.15.0", "status": "waiting_for_account"}, + capsys) + assert "0.15.0" in out and "warning" in out and "not replaced" in out + + +# ── a graceful stop, through the control API ─────────────────────────────── + + +def test_the_control_api_asks_the_daemon_to_stop(): + from fastapi.testclient import TestClient + from meshbay_node.ui import create_ui_app + + asked = [] + app = create_ui_app({"ui_token": "tok", "request_shutdown": lambda: asked.append(1)}) + client = TestClient(app) + assert client.post("/api/shutdown").status_code in (401, 403), "behind the token" + assert asked == [] + r = client.post("/api/shutdown?t=tok") + assert r.status_code == 200 and r.json() == {"stopping": True} + assert asked == [1] + + +def test_a_stop_request_ends_the_wait_for_the_hub(): + """A node waiting for its account is the one people stop and restart.""" + import asyncio + + from meshbay_node.daemon import NodeDaemon + + async def go(): + daemon = NodeDaemon.__new__(NodeDaemon) + daemon._stop_event = asyncio.Event() + asyncio.get_running_loop().call_later(0.1, daemon._stop_event.set) + t0 = time.monotonic() + stopped = await daemon._pause(30) + return stopped, time.monotonic() - t0 + + stopped, took = asyncio.run(go()) + assert stopped and took < 5 + + +class _RealNode: + """A real daemon process in a throwaway profile, pointed at a hub that is + not there (it waits for it, control API up).""" + + def __init__(self, tmp_path): + import socket + + with socket.socket() as s: + s.bind(("127.0.0.1", 0)) + self.port = s.getsockname()[1] + self.tmp = tmp_path + home = tmp_path / "home" + self.conf = tmp_path / "node.toml" + unlock = tmp_path / "unlock.key" + unlock.write_text("x" * 43, encoding="utf-8") + self.data = tmp_path / "data" + self.conf.write_text( + f'data_dir = "{self.data.as_posix()}"\n' + '[hub]\nurl = "http://127.0.0.1:1"\nusername = "probe"\n' + f'[node]\nquic_enabled = false\nui_port = {self.port}\n' + f'[keystore]\npath = "{(tmp_path / "keystore.enc").as_posix()}"\n' + f'unlock_file = "{unlock.as_posix()}"\n', encoding="utf-8") + self.env = {**os.environ, "HOME": str(home), "LOCALAPPDATA": str(home), + "USERPROFILE": str(home)} + self.procs = [] + + def spawn(self, name: str): + import subprocess + import sys + + out = self.tmp / f"{name}.txt" + with open(out, "w", encoding="utf-8") as f: + proc = subprocess.Popen( + [sys.executable, "-c", "from meshbay_node.daemon import main; main()", + "--config", str(self.conf)], + env=self.env, stderr=f, stdout=f) + proc.out = out + self.procs.append(proc) + return proc + + def wait_up(self, proc) -> None: + deadline = time.monotonic() + 40 + while not (self.data / "ui-token").exists() and time.monotonic() < deadline: + assert proc.poll() is None, proc.out.read_text(encoding="utf-8") + time.sleep(0.2) + time.sleep(1) + + def kill_all(self) -> None: + for p in self.procs: + if p.poll() is None: + p.kill() + + +def test_a_real_daemon_shuts_down_properly_when_asked(tmp_path): + """The whole path, on a real process: before this, nine stops of a Windows + node in a row logged not one shutdown -- every one was a TerminateProcess.""" + node = _RealNode(tmp_path) + proc = node.spawn("first") + try: + node.wait_up(proc) + assert plat.request_graceful_stop(node.data, node.port, timeout=15) + assert proc.wait(timeout=10) is not None + err = proc.out.read_text(encoding="utf-8") + assert "Node stopped" in err, err[-2000:] + finally: + node.kill_all() + + +def test_a_second_node_leaves_the_running_one_alone_and_says_why(tmp_path): + """Found on a real install: the sign-in launcher run while a node was up + started a second one, which wrote its token over the first's, failed to bind + inside uvicorn's task and exited with nothing in the log. The node that kept + running then refused every stop, status and restart -- the token file named + a process that no longer existed.""" + node = _RealNode(tmp_path) + first = node.spawn("first") + try: + node.wait_up(first) + token = (node.data / "ui-token").read_text(encoding="utf-8") + + second = node.spawn("second") + assert second.wait(timeout=40) == 2 + said = second.out.read_text(encoding="utf-8") + assert f"127.0.0.1:{node.port} is taken" in said, said[-2000:] + assert "already running" in said + + assert (node.data / "ui-token").read_text(encoding="utf-8") == token + assert first.poll() is None + assert plat.request_graceful_stop(node.data, node.port, timeout=15), ( + "the running node can no longer be asked to stop") + finally: + node.kill_all() + + +def test_the_control_port_is_exclusive(): + """On Windows SO_REUSEADDR would let a second socket share a port in use.""" + from meshbay_node.daemon import ControlPortTaken, bind_control_port + + first = bind_control_port(0) + try: + port = first.getsockname()[1] + with pytest.raises(ControlPortTaken, match="already running"): + bind_control_port(port) + finally: + first.close() + + +# ── "No operator paired" only when the roster says so ────────────────────── + + +async def test_a_node_that_has_not_read_its_roster_does_not_claim_there_is_no_operator(): + from meshbay_node import ops + result = await ops.list_groups({"config": None, "groups_ctx": {}}) + assert result["operator_paired"] is None + + +async def test_a_node_with_a_roster_reports_the_operator(tmp_path): + from meshbay_node import ops + from meshbay_node.roster import Roster + + roster = Roster(db_path=tmp_path / "roster.db") + await roster.open() + try: + result = await ops.list_groups({"config": None, "groups_ctx": {}, "roster": roster}) + finally: + await roster.close() + assert result["operator_paired"] is False |