aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-27 22:20:35 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-27 22:20:35 +0200
commit7662484cae8e74b7d9aa383bd6cd0dad4690aadc (patch)
tree057609311e339c839dcad04b33de62e2541a9943 /packages/meshbay-node/src/meshbay_node
parent326d7796c79f616f5a0f2058386df2bded657c78 (diff)
downloadmeshbay-7662484cae8e74b7d9aa383bd6cd0dad4690aadc.tar.gz
fix(node): a Windows daemon that stops properly, starts honestly and runs once
Found by installing the builds and driving every startup mode live: - Stop through the node's own control API first (POST /api/shutdown, loopback and per-run token): the one channel that reaches a daemon in any session without elevation -- a service node runs in session 0 -- and the one that runs its shutdown. Then Task Scheduler, then a forced stop. Nine stops in a row used to log no shutdown at all: each was a TerminateProcess. - The forced stop spares the command running it. The frozen meshbay-node.exe is the daemon and every CLI verb, so `taskkill /IM meshbay-node.exe` killed `autostart stop` and `restart-daemon` themselves: exit 1, no output, and no node after a restart. It excludes its own pid and its parent's, and /T takes a venv launcher's python child and a daemon's ffmpeg children with it. - Start and restart report the version that answered, never "started" about a node nobody asked; `service start` says so when no node answered, and where the log is. - A second instance fails before it touches anything. The daemon wrote ui-token, then failed to bind inside uvicorn's task and exited with the reason on a hidden console; the node still running then refused every stop and status, its token file naming a dead process. The control port is now bound first (exclusively on Windows, where SO_REUSEADDR would share it), and a refusal is logged and exits 2. Linux had the same order. - The daemon logs to %LOCALAPPDATA%\meshbay\state\node.log: Task Scheduler discards its stderr. Only the daemon run opens it, never a CLI verb. - Hub sign-in waits are interruptible, a stop requested before the node is up is honoured, and a hub that answers 429 or restarts leaves the node in waiting_for_hub rather than looking dead. - operator_paired is null until the roster is read, instead of a false that showed "No operator paired" about a node whose pairing was intact. The node test conftest also points HOME, USERPROFILE, LOCALAPPDATA and APPDATA at a throwaway directory for every test, and keeps log_file() away from the developer's own node: redirecting HOME alone isolates nothing on Windows, and the CLI tests had been writing invite and pairing codes into the real profile. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node')
-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
6 files changed, 392 insertions, 125 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")