diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/daemon.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 90 |
1 files changed, 61 insertions, 29 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 220a908..45a6b9a 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -44,6 +44,7 @@ from meshbay_common.protocol import MNP from meshbay_node import uploads as uploads_mod from meshbay_node.audit import RETENTION_DAYS as AUDIT_RETENTION_DAYS from meshbay_node.audit import AuditStore +from meshbay_node.blocklist import ContentBlocklist from meshbay_node.bundle_store import BundleStore from meshbay_node.chat.store import ChatStore from meshbay_node.cli.dispatch import run, start @@ -114,6 +115,10 @@ class NodeDaemon(EnrichmentMixin): # Persisted so a restart does not silently un-revoke everyone (H4) self._denylist = ( Denylist(path=config.data_dir / "denylist.json") if Denylist else None) + # The hub's content blocklist, applied in public groups only (§7.5). + # Persisted for the same reason: a restart while the hub is unreachable + # must not serve again what had stopped being served. + self._blocklist = ContentBlocklist(config.data_dir / "blocklist.json") self._chat_stores: dict[str, ChatStore] = {} # One instance, shared by every group's DirectoryIndexer — see # indexer/cache.py's docstring for why this stopped being per-group. @@ -525,6 +530,7 @@ class NodeDaemon(EnrichmentMixin): self._webrtc._ctx["pk_x25519_b64"] = keys.pk_x25519_b64 self._webrtc._ctx["roster"] = self._roster + self._webrtc._ctx["blocklist"] = self._blocklist # The MNP adapter calls the same operations as the loopback API # (meshbay_node.ops), and those take the daemon's state. Handing # the transport a second set of lookups is how two paths to one @@ -611,11 +617,16 @@ class NodeDaemon(EnrichmentMixin): except Exception as e: log.warning("Invalid revocation token: %s", e) + async def on_connected(): + await self._sync_blocklist(hub) + ws_task = asyncio.create_task(hub.maintain_ws( on_incoming=on_incoming, on_revocation=on_revocation, on_webrtc_offer=on_webrtc_offer, group_ids=lambda: list((self._state.get("groups_ctx") or {}).keys()), + on_connected=on_connected, + on_blocklist=self._on_blocklist_update, )) self._tasks.append(ws_task) log.info("Hub WS task started") @@ -667,12 +678,6 @@ class NodeDaemon(EnrichmentMixin): await indexer.initial_scan() log.info("Background scan complete for %s: %d files", name, indexer.index.count) - # Swarm registration for public groups (after files are known). - if gctx.get("visibility") == "public": - endpoint = f"webrtc:{self._config.node.quic_port}" - hashes = [e.id for e in gctx["index"].entries] - if hashes: - await self._register_swarm(hashes, endpoint) # initial_scan() itself never calls on_change (it predates # the concept — every existing caller only cared about the # scan finishing, not about notifying anyone) — but Videos @@ -1454,9 +1459,10 @@ class NodeDaemon(EnrichmentMixin): # node refuses every handshake while the GEK is None (NS8) — so this is # "nobody is listening", not a case to send in clear for. if peers and idx.gek: - msg = (index_delta_message(idx, delta, indexer.roots) + hidden = self._blocklist if self._is_public(group_id) else () + msg = (index_delta_message(idx, delta, indexer.roots, hidden) if delta is not None - else index_sync_message(idx, indexer.roots)) + else index_sync_message(idx, indexer.roots, hidden)) pushed = 0 for session in peers: try: @@ -1468,20 +1474,53 @@ class NodeDaemon(EnrichmentMixin): log.info("Index %s pushed to %d WebRTC peers", "delta" if delta is not None else "sync", pushed) - # 11.9 — Register file hashes with hub swarm table (public groups only, H7) - group_cfg = next( - (g for g in self._config.groups if g.id == group_id), None) - if (self._hub and self._state.get("endpoint_hint") - and group_cfg and group_cfg.visibility == "public"): - # Only the newly added hashes once there is a delta to know them - # from — registering the whole library again on every change is - # the same O(changes x library size) cost the delta above exists - # to avoid. - hashes = ([e.id for e in delta.additions] if delta is not None - else [e.id for e in idx.entries]) - if hashes: - endpoint = f"webrtc:{self._config.node.quic_port}" - spawn(self._register_swarm(hashes, endpoint)) + # ── Content blocklist (public groups, §7.5) ────────────────────────────── + + def _is_public(self, group_id: str) -> bool: + return any(g.id == group_id and g.visibility == "public" + for g in self._config.groups) + + def _hosts_public_group(self) -> bool: + return any(g.visibility == "public" for g in self._config.groups) + + async def _sync_blocklist(self, hub) -> None: + """The hub's whole list, on every connection. Only a node hosting a + public group asks: nothing else here could have been blocked.""" + if not self._hosts_public_group(): + return + try: + hashes = await hub.fetch_blocklist() + except Exception as e: + # Keep applying the list on disk; the next connection tries again. + log.warning("Content blocklist sync failed: %s", e) + return + if self._blocklist.replace(hashes): + log.info("Content blocklist: %d hashes", len(self._blocklist)) + self._push_public_indexes() + + def _on_blocklist_update(self, add: list, remove: list) -> None: + if not self._hosts_public_group(): + return + if self._blocklist.apply(add, remove): + log.info("Content blocklist updated: +%d -%d", len(add), len(remove)) + self._push_public_indexes() + + def _push_public_indexes(self) -> None: + """Resend each public group's whole index, so what the list now hides + leaves every connected member's view, and what it released comes back.""" + if not self._webrtc: + return + for indexer in self._indexers: + idx = indexer.index + if not self._is_public(idx.group_id) or not idx.gek: + continue + msg = index_sync_message(idx, indexer.roots, self._blocklist) + for session in list(self._webrtc._sessions.values()): + if session._group_id == idx.group_id: + try: + session._send(msg) + except Exception: + pass def _drop_group_sessions(self, group_id: str) -> None: """Close live sessions for a revoked group (H4).""" @@ -1492,13 +1531,6 @@ class NodeDaemon(EnrichmentMixin): spawn(session.close()) log.info("Dropped session for revoked group %s", group_id[:8]) - 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" |