aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-node/src/meshbay_node/cli/dispatch.py27
-rw-r--r--packages/meshbay-node/src/meshbay_node/cli/lifecycle.py154
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py124
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/groups.py5
-rw-r--r--packages/meshbay-node/src/meshbay_node/platform.py188
-rw-r--r--packages/meshbay-node/src/meshbay_node/ui/app.py19
-rw-r--r--packages/meshbay-node/tests/conftest.py42
-rw-r--r--packages/meshbay-node/tests/test_login_retry_is_resilient.py20
-rw-r--r--packages/meshbay-node/tests/test_platform.py252
-rw-r--r--packages/meshbay-node/tests/test_windows_service_diagnosis.py332
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