""" MeshBay Node daemon — main process. Startup sequence: 1. Load config (~/.config/meshbay/node.toml) 2. Load or create keystore (Argon2id unlock) 3. Connect to hub: register → login → announce node 4. Fetch GEK bundle from hub (if group configured) 5. Start directory indexer (watchdog) 6. Create chat stores (one SQLite DB per group) 7. Create WebRTC transport (browser clients via DataChannel) 8. Start QUIC+TCP chunk servers (native clients) 9. Start HTTP file API (public content) 10. Start hub WebSocket (signaling, revocations, WebRTC offers) 11. Start local web UI on node.ui_port (localhost only) 12. Run until SIGINT/SIGTERM Usage: meshbay-node # interactive password prompt meshbay-node --config /path # custom config meshbay-node init # write example config + create keystore meshbay-node --calibrate-argon2 # benchmark Argon2id, suggest parameters """ import asyncio import base64 import json import logging import signal import sys from pathlib import Path import uvicorn from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey from meshbay_common import MNP_VERSION from meshbay_common.protocol import MNP from meshbay_node.audit import AuditStore from meshbay_node.bundle_store import BundleStore from meshbay_node.chat.store import ChatStore from meshbay_node.config import Config, DEFAULT_CONFIG_PATH, load_config, write_example_config from meshbay_node.hub_client import HubClient, HubConfig from meshbay_node.indexer import DirectoryIndexer from meshbay_node.keystore import NodeKeys, load_or_create_keystore from meshbay_node.transport import ( ChunkServer, Denylist, QUIC_AVAILABLE, WEBRTC_AVAILABLE, create_http_app, ) if QUIC_AVAILABLE: from meshbay_node.transport import QuicChunkServer if WEBRTC_AVAILABLE: from meshbay_node.transport import WebRTCTransport log = logging.getLogger(__name__) # ── Argon2id calibration ────────────────────────────────────────────────────── def calibrate_argon2(target_ms: int = 500) -> None: """Benchmark Argon2id and suggest parameters targeting ~target_ms.""" import time import os print(f"Calibrating Argon2id (target: {target_ms}ms) ...") salt = os.urandom(16) for mem in [65536, 131072, 262144, 524288]: from cryptography.hazmat.primitives.kdf.argon2 import Argon2id t0 = time.perf_counter() Argon2id(salt=salt, length=32, iterations=3, lanes=4, memory_cost=mem).derive(b"benchmark") elapsed_ms = (time.perf_counter() - t0) * 1000 print(f" memory_cost={mem:>7} ({mem//1024:>4}MB): {elapsed_ms:.0f}ms", end="") if abs(elapsed_ms - target_ms) < target_ms * 0.3: print(" ← recommended") else: print() print("Set memory_cost in meshbay_common/crypto.py: ARGON2_MEMORY_COST") # ── Hub WS sender bridge ───────────────────────────────────────────────────── class _WsSender: """Thin bridge so WebRTC context can call hub_ws.send() for chat_notify.""" def __init__(self, hub_client: HubClient): self._hub = hub_client async def send(self, data: str) -> None: await self._hub.send_ws(data) # ── Daemon ──────────────────────────────────────────────────────────────────── class NodeDaemon: def __init__(self, config: Config): self._config = config self._state: dict = { "status": "starting", "hub_url": config.hub.url, "username": config.hub.username, "groups": [g.name for g in config.groups], "node_port": config.node.port, "quic_port": config.node.quic_port, "endpoint_hint": None, "indexes": {}, } self._tcp_server: ChunkServer | None = None self._quic_server = None self._webrtc = None self._denylist = Denylist() if Denylist else None self._chat_stores: dict[str, ChatStore] = {} self._audit_store: AuditStore | None = None self._bundle_store: BundleStore | None = None self._indexers: list[DirectoryIndexer] = [] self._tasks: list[asyncio.Task] = [] self._hub: HubClient | None = None self._http_servers: list[uvicorn.Server] = [] async def run(self) -> None: log.info("MeshBay Node starting up") # 1. Keystore keys = load_or_create_keystore( path=self._config.keystore.path, unlock_file=self._config.keystore.unlock_file, ) log.info("Keys loaded: %s", keys.pk_ed25519_b64[:16]) # 2. Start admin UI early (so operator can copy node key before hub login) self._state["pk_node_ed25519"] = keys.pk_ed25519_b64 self._state["config"] = self._config from meshbay_node.ui import create_ui_app ui_app = create_ui_app(self._state) ui_cfg = uvicorn.Config( ui_app, host="127.0.0.1", port=self._config.node.ui_port, log_level="warning", ) ui_server = uvicorn.Server(ui_cfg) self._tasks.append(asyncio.create_task(ui_server.serve())) log.info("Admin UI at http://localhost:%d", self._config.node.ui_port) # 3. Hub connection (Ed25519 auth — retries until node key is linked) hub_cfg = HubConfig( hub_url=self._config.hub.url, username=self._config.hub.username, ) async with HubClient(hub_cfg, keys) as hub: self._hub = hub session = await self._login_with_retry(hub) self._state["endpoint_hint"] = session.node_id # 4. Bundle store (P2P GEK bundles) data_dir = self._config.data_dir data_dir.mkdir(parents=True, exist_ok=True) self._bundle_store = BundleStore(db_path=data_dir / "bundles.db") await self._bundle_store.open() log.info("Bundle store opened: %s", data_dir / "bundles.db") # X25519 key material for GEK unwrapping from cryptography.hazmat.primitives import serialization sk_x_raw = keys.sk_x25519.private_bytes( serialization.Encoding.Raw, serialization.PrivateFormat.Raw, serialization.NoEncryption()) pk_x_raw = base64.b64decode(keys.pk_x25519_b64) # 4. Build per-group contexts groups_ctx: dict[str, dict] = {} for group_cfg in self._config.groups: if not group_cfg.id or not group_cfg.shared_dir: log.warning("Group %r missing id or shared_dir — skipping", group_cfg.name) continue shared_root = Path(group_cfg.shared_dir).expanduser().resolve() if not shared_root.exists(): log.warning("Shared dir not found: %s — skipping group %s", shared_root, group_cfg.name) continue gek = None if group_cfg.visibility == "private": gek = await self._load_gek( group_cfg.id, session.user_id, sk_x_raw, pk_x_raw) if gek: log.info("GEK loaded for group %s", group_cfg.id[:8]) else: log.info("No GEK yet for group %s — will accept first setup", group_cfg.name) indexer = DirectoryIndexer( root=shared_root, group_id=group_cfg.id, sk_node=keys.sk_ed25519, gek=gek, on_change=self._on_index_change, ) await indexer.start() self._indexers.append(indexer) self._state["indexes"][group_cfg.id] = indexer.index log.info("Indexing group %s: %s (%d files)", group_cfg.name, shared_root, indexer.index.count) groups_ctx[group_cfg.id] = { "gek": gek, "shared_root": shared_root, "index": indexer.index, } if not groups_ctx: log.error("No valid groups configured — exiting") return # 5. Chat stores (one SQLite DB per group) for gid in groups_ctx: chat_db = data_dir / gid[:16] / "chat.db" store = ChatStore(db_path=chat_db) await store.open() self._chat_stores[gid] = store groups_ctx[gid]["chat_store"] = store log.info("Chat stores opened: %d groups", len(self._chat_stores)) # 6. Audit store (legal compliance — IP + action logging) audit_db = data_dir / "audit.db" self._audit_store = AuditStore(db_path=audit_db) await self._audit_store.open() log.info("Audit store opened: %s", audit_db) # 5. Denylist denylist = self._denylist # 6. WebRTC transport (browser clients) first = next(iter(groups_ctx.values())) if WEBRTC_AVAILABLE: self._webrtc = WebRTCTransport( sk_node=keys.sk_ed25519, hub_pk_pem=session.hub_pk_pem, gek=first["gek"], shared_root=first["shared_root"], index=first["index"], groups=groups_ctx, denylist=denylist, ) self._webrtc._ctx["chat_store"] = first.get("chat_store") self._webrtc._ctx["hub_ws"] = _WsSender(hub) self._webrtc._ctx["node_user_id"] = session.user_id self._webrtc._ctx["audit_store"] = self._audit_store self._webrtc._ctx["bundle_store"] = self._bundle_store self._webrtc._ctx["sk_x25519_raw"] = sk_x_raw self._webrtc._ctx["pk_x25519_raw"] = pk_x_raw self._webrtc._ctx["pk_x25519_b64"] = keys.pk_x25519_b64 admin_pk = self._resolve_admin_pk(keys) if admin_pk: self._webrtc._ctx["admin_pk_ed25519"] = admin_pk log.info("Admin Ed25519 key pinned for node sovereignty") else: log.warning("No admin_pk_ed25519 — admin operations disabled") log.info("WebRTC transport ready") else: log.warning("WebRTC not available (aiortc not installed)") # 7. QUIC + TCP chunk servers if QUIC_AVAILABLE: self._quic_server = QuicChunkServer( sk_node=keys.sk_ed25519, hub_pk_pem=session.hub_pk_pem, gek=first["gek"], shared_root=first["shared_root"], index=first["index"], host="::", port=self._config.node.quic_port, groups=groups_ctx, denylist=denylist, ) await self._quic_server.start() log.info("QUIC server on port %d (%d groups)", self._config.node.quic_port, len(groups_ctx)) self._tcp_server = ChunkServer( sk_node=keys.sk_ed25519, hub_pk_pem=session.hub_pk_pem, gek=first["gek"], shared_root=first["shared_root"], index=first["index"], host="0.0.0.0", port=self._config.node.port, groups=groups_ctx, ) await self._tcp_server.start() log.info("TCP+TLS server on port %d", self._config.node.port) # 8. Hub WebSocket (signaling + revocations + WebRTC offers) async def on_webrtc_offer(sdp, peer_id, ice_candidates): if not self._webrtc: return None try: answer_sdp, answer_ice = await self._webrtc.handle_offer( sdp, peer_id) log.info("WebRTC answer for peer=%s (%d peers)", peer_id, self._webrtc.active_peers) return (answer_sdp, answer_ice) except Exception as e: log.error("WebRTC offer failed: %s", e) return None async def on_incoming(peer_ip, peer_port): if self._quic_server: self._quic_server.punch_nat(peer_ip, peer_port) def on_revocation(token): if denylist and token: import jwt as _jwt try: payload = _jwt.decode( token, session.hub_pk_pem, algorithms=["EdDSA"], options={"verify_exp": False}) target = payload.get("target") tid = payload.get("target_id", "") if target == "user": denylist.deny_user(tid) elif target == "jti": denylist.deny_jti(tid) except Exception as e: log.warning("Invalid revocation token: %s", e) ws_task = asyncio.create_task(hub.maintain_ws( on_incoming=on_incoming, on_revocation=on_revocation, on_webrtc_offer=on_webrtc_offer, group_ids=list(groups_ctx.keys()), )) self._tasks.append(ws_task) log.info("Hub WS task started") # 9. HTTP file API (one per group) for gid, gctx in groups_ctx.items(): group_cfg = next( (g for g in self._config.groups if g.id == gid), None) if not group_cfg: continue http_app = create_http_app( sk_node=keys.sk_ed25519, hub_pk_pem=session.hub_pk_pem, shared_root=gctx["shared_root"], index=gctx["index"], group_id=gid, group_name=group_cfg.name, gek=gctx.get("gek"), ) http_cfg = uvicorn.Config( http_app, host="0.0.0.0", port=group_cfg.http_port, log_level="warning", ) http_server = uvicorn.Server(http_cfg) self._http_servers.append(http_server) self._tasks.append(asyncio.create_task(http_server.serve())) log.info("HTTP API on port %d for group %s", group_cfg.http_port, group_cfg.name) # 10. Update admin UI state (UI already running from step 2) self._state["groups_ctx"] = groups_ctx self._state["audit_store"] = self._audit_store self._state["bundle_store"] = self._bundle_store self._state["webrtc"] = self._webrtc self._state["hub"] = hub self._state["pk_x25519_raw"] = pk_x_raw self._state["status"] = "running" log.info("Node ready — %d groups, WebRTC=%s, QUIC=%s", len(groups_ctx), "yes" if self._webrtc else "no", "yes" if self._quic_server else "no") # 11. Initial swarm registration endpoint = f"webrtc:{self._config.node.quic_port}" for gctx in groups_ctx.values(): hashes = [e.id for e in gctx["index"].entries] if hashes: asyncio.ensure_future(self._register_swarm(hashes, endpoint)) # 12. Wait for shutdown stop_event = asyncio.Event() loop = asyncio.get_event_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, stop_event.set) await stop_event.wait() await self._shutdown() async def _login_with_retry(self, hub: HubClient): """Login to hub, retrying if the node key hasn't been linked yet.""" import httpx as _httpx while True: try: return await hub.startup(endpoint_hint=None) except _httpx.HTTPStatusError as e: body = e.response.text if hasattr(e.response, 'text') else '' if e.response.status_code == 401 and "No node key" in body: self._state["status"] = "waiting_for_node_key" log.warning( "Node key not linked — open admin UI at " "http://localhost:%d, copy the key, and paste it in " "Settings > Link Node on the hub. Retrying in 30s...", self._config.node.ui_port, ) await asyncio.sleep(30) else: raise except Exception as e: log.warning("Hub login failed: %s — retrying in 10s", e) await asyncio.sleep(10) async def _load_gek( self, group_id: str, node_user_id: str, sk_x_raw: bytes, pk_x_raw: bytes, ) -> bytes | None: """Load GEK from local bundle store (node-only, hub never touches crypto).""" from meshbay_common.crypto import unwrap_gek_aes if not self._bundle_store: return None # Try node-specific bundle first (stored by init_gek for daemon reload), # then fall back to operator's user bundle (legacy / pre-dual-key) for user_key in [f"_node_{node_user_id}", node_user_id]: bundle = await self._bundle_store.fetch(group_id, user_key) if not bundle: continue try: gek = unwrap_gek_aes(bundle, sk_x_raw, pk_x_raw) log.info("GEK loaded from local bundle store for group %s (key=%s)", group_id[:8], user_key[:16]) return gek except Exception as e: log.debug("Failed to unwrap GEK bundle (key=%s): %s", user_key[:16], e) log.warning("No unwrappable GEK bundle found for group %s", group_id[:8]) return None def _resolve_admin_pk(self, keys: NodeKeys) -> Ed25519PublicKey | None: """Resolve the admin Ed25519 public key: config → auto-pin from node keystore.""" if self._config.admin_pk_ed25519: try: raw = base64.b64decode(self._config.admin_pk_ed25519) return Ed25519PublicKey.from_public_bytes(raw) except Exception as e: log.error("Invalid admin_pk_ed25519 in config: %s", e) return None pk = keys.sk_ed25519.public_key() from meshbay_common.crypto import pk_to_b64 pk_b64 = pk_to_b64(pk) log.info("Auto-pinning admin key from node keystore: %s", pk_b64[:16]) return pk async def _on_index_change(self, indexer: DirectoryIndexer) -> None: """Called when a DirectoryIndexer detects file changes.""" group_id = indexer.group_id idx = indexer.index log.info("Index changed for group %s: %d files (v%d)", group_id[:8], idx.count, idx.version) # 11.5 — Push updated index to connected WebRTC peers in this group if self._webrtc: entries = [ { "id": e.id, "name": e.name, "path": e.path, "size": e.size, "type": e.type, "added_at": e.added_at, } for e in idx.entries ] sync_msg = { "type": MNP.INDEX_SYNC, "v": MNP_VERSION, "group_id": idx.group_id, "version": idx.version, "entries": entries, } pushed = 0 for session in list(self._webrtc._sessions.values()): if session._group_id == group_id: try: session._send(sync_msg) pushed += 1 except Exception: pass if pushed: log.info("Index pushed to %d WebRTC peers", pushed) # 11.9 — Register file hashes with hub swarm table if self._hub and self._state.get("endpoint_hint"): hashes = [e.id for e in idx.entries] if hashes: endpoint = f"webrtc:{self._config.node.quic_port}" asyncio.ensure_future(self._register_swarm(hashes, endpoint)) async def _register_swarm(self, hashes: list[str], endpoint: str) -> None: try: n = await self._hub.register_swarm(hashes, endpoint) log.info("Swarm: registered %d/%d hashes", n, len(hashes)) except Exception as e: log.warning("Swarm registration failed: %s", e) async def _shutdown(self) -> None: log.info("Shutting down...") self._state["status"] = "stopping" for task in self._tasks: task.cancel() for task in self._tasks: try: await task except (asyncio.CancelledError, Exception): pass if self._webrtc: await self._webrtc.close_all() if self._audit_store: await self._audit_store.close() if self._bundle_store: await self._bundle_store.close() for store in self._chat_stores.values(): await store.close() for indexer in self._indexers: await indexer.stop() if self._quic_server: await self._quic_server.stop() if self._tcp_server: await self._tcp_server.stop() for server in self._http_servers: server.should_exit = True log.info("Node stopped") # ── Entry point ─────────────────────────────────────────────────────────────── def main() -> None: import argparse parser = argparse.ArgumentParser(description="MeshBay Node daemon") parser.add_argument("command", nargs="?", choices=["init", "calibrate-argon2"], help="init: write example config | calibrate-argon2: benchmark") parser.add_argument("--config", type=Path, default=None, help="Config file path") parser.add_argument("--log-level", default="INFO", choices=["DEBUG", "INFO", "WARNING", "ERROR"]) args = parser.parse_args() logging.basicConfig( level=getattr(logging, args.log_level), format="%(asctime)s %(levelname)-8s %(name)s: %(message)s", ) if args.command == "init": write_example_config() print("Example config written. Edit it and run: meshbay-node") return if args.command == "calibrate-argon2": calibrate_argon2() return 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") sys.exit(1) daemon = NodeDaemon(cfg) asyncio.run(daemon.run()) if __name__ == "__main__": main()