diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/daemon.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 111 |
1 files changed, 71 insertions, 40 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index bdac6b8..93ba3c4 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -31,6 +31,7 @@ from meshbay_node.hub_client import HubClient, HubConfig from meshbay_node.indexer import DirectoryIndexer from meshbay_node.keystore import load_or_create_keystore from meshbay_node.transport import ChunkServer +from meshbay_node.transport.quic_server import QuicChunkServer from meshbay_node.ui import create_ui_app log = logging.getLogger(__name__) @@ -69,13 +70,14 @@ class NodeDaemon: "status": "starting", "hub_url": config.hub.url, "username": config.hub.username, - "group_id": config.group.id, - "group_name": config.group.name, + "groups": [g.name for g in config.groups], "node_port": config.node.port, + "quic_port": config.node.quic_port, "endpoint_hint": None, - "index": None, + "indexes": {}, } - self._server: ChunkServer | None = None + self._tcp_server: ChunkServer | None = None + self._quic_server: QuicChunkServer | None = None self._indexers: list[DirectoryIndexer] = [] self._tasks: list[asyncio.Task] = [] @@ -99,52 +101,79 @@ class NodeDaemon: session = await hub.startup(endpoint_hint=None) self._state["endpoint_hint"] = session.node_id - # 3. Fetch GEK if group configured - if self._config.group.id: - try: - gek = await hub.fetch_gek(self._config.group.id) - keys.gek = gek - log.info("GEK loaded for group %s", self._config.group.id[:8]) - except LookupError: - log.warning("No GEK bundle found for group %s — " - "wait for admin to add you", self._config.group.id[:8]) - - # 4. Directory indexers - async def on_index_change(indexer: DirectoryIndexer) -> None: - self._state["index"] = indexer.index + # 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 - for shared_dir in self._config.node.shared_dirs: - d = Path(shared_dir).expanduser().resolve() - if not d.exists(): - log.warning("Shared directory not found: %s — skipping", d) + 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=d, - group_id=self._config.group.id, + root=shared_root, + group_id=group_cfg.id, sk_node=keys.sk_ed25519, - gek=keys.gek, - on_change=on_index_change, + gek=gek, ) await indexer.start() self._indexers.append(indexer) - self._state["index"] = indexer.index - log.info("Indexing: %s (%d files)", d, indexer.index.count) + 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, + } + + # 4. QUIC chunk server (primary transport, all groups on one port) + if groups_ctx: + first = next(iter(groups_ctx.values())) + 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, + ) + await self._quic_server.start() + log.info("QUIC server on port %d (%d groups)", + self._config.node.quic_port, len(groups_ctx)) - # 5. Chunk server - if self._indexers and keys.gek: - self._server = ChunkServer( + # TCP+TLS server (fallback transport, same groups) + self._tcp_server = ChunkServer( sk_node=keys.sk_ed25519, hub_pk_pem=session.hub_pk_pem, - gek=keys.gek, - shared_root=Path(self._config.node.shared_dirs[0]).expanduser(), - index=self._indexers[0].index, + 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._server.start() - log.info("Chunk server on port %d", self._config.node.port) + await self._tcp_server.start() + log.info("TCP+TLS server on port %d", self._config.node.port) - # 6. Local web UI + # 5. Local web UI ui_app = create_ui_app(self._state) ui_cfg = uvicorn.Config( ui_app, @@ -157,9 +186,9 @@ class NodeDaemon: log.info("Local UI at http://localhost:%d", self._config.node.ui_port) self._state["status"] = "running" - log.info("Node ready") + log.info("Node ready — %d groups", len(groups_ctx)) - # 7. Wait for shutdown + # 6. Wait for shutdown stop_event = asyncio.Event() loop = asyncio.get_event_loop() for sig in (signal.SIGINT, signal.SIGTERM): @@ -176,8 +205,10 @@ class NodeDaemon: task.cancel() for indexer in self._indexers: await indexer.stop() - if self._server: - await self._server.stop() + if self._quic_server: + await self._quic_server.stop() + if self._tcp_server: + await self._tcp_server.stop() log.info("Node stopped") |