""" 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 json import logging import signal import sys from pathlib import Path import uvicorn from meshbay_common import MNP_VERSION from meshbay_common.protocol import MNP from meshbay_node.audit import AuditStore from meshbay_node.chat.store import ChatStore from meshbay_node.config import Config, 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._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. Hub connection hub_cfg = HubConfig( hub_url=self._config.hub.url, username=self._config.hub.username, password=self._config.hub.password, ) async with HubClient(hub_cfg, keys) as hub: self._hub = hub session = await hub.startup(endpoint_hint=None) self._state["endpoint_hint"] = session.node_id # 3. 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": try: gek = await hub.fetch_gek(group_cfg.id) log.info("GEK loaded for group %s", group_cfg.id[:8]) except LookupError: log.warning("No GEK for group %s — skipping", group_cfg.name) continue 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 # 4. Chat stores (one SQLite DB per group) data_dir = self._config.data_dir data_dir.mkdir(parents=True, exist_ok=True) 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)) # 4b. 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 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. Local web UI self._state["groups_ctx"] = groups_ctx self._state["config"] = self._config self._state["audit_store"] = self._audit_store self._state["webrtc"] = self._webrtc self._state["hub"] = hub 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("Local UI at http://localhost:%d", self._config.node.ui_port) 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 _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() 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) 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()