From 07480eb3f8ad0bb4369ac8c41df7c4140b108d0e Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Mon, 28 Sep 2026 21:49:14 +0200 Subject: feat: nodes apply the content blocklist in their public groups A node hosting a public group syncs the hub's blocklist on every connection (paged, node token only) and applies pushed changes. A blocked file leaves the index and is refused (content_blocked); private groups are untouched. The unused per-hash check route is gone. Co-Authored-By: Claude Opus 5.5 --- .../meshbay-node/src/meshbay_node/blocklist.py | 83 ++++++++++++++++++++++ packages/meshbay-node/src/meshbay_node/daemon.py | 64 ++++++++++++++++- .../meshbay-node/src/meshbay_node/hub_client.py | 31 ++++++++ .../meshbay_node/transport/webrtc/apps/music.py | 2 + .../transport/webrtc/apps/streaming.py | 2 + .../transport/webrtc/apps/subtitles.py | 2 + .../src/meshbay_node/transport/webrtc/core.py | 19 +++++ .../src/meshbay_node/transport/webrtc/files.py | 12 +++- .../src/meshbay_node/transport/wire.py | 16 +++-- 9 files changed, 223 insertions(+), 8 deletions(-) create mode 100644 packages/meshbay-node/src/meshbay_node/blocklist.py (limited to 'packages/meshbay-node/src') diff --git a/packages/meshbay-node/src/meshbay_node/blocklist.py b/packages/meshbay-node/src/meshbay_node/blocklist.py new file mode 100644 index 0000000..29d8666 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/blocklist.py @@ -0,0 +1,83 @@ +""" +The hub's content blocklist, as this node applies it (docs/MESHBAY_DESIGN.md §7.5). + +Hashes of public content that the hub's moderation has blocked. A node hosting a +public group stops serving those files there: they leave the index members are +sent, and a request for one is refused. Nothing is deleted from the operator's +disk — the node stops serving, and what the operator keeps is theirs to decide. + +Private groups are untouched. The hub never learns what a private group holds, +so there is nothing it could have blocked in one, and nothing here reaches them. + +Kept on disk as well as in memory, so a node that restarts while the hub is +unreachable does not serve again, meanwhile, what it had already stopped serving. +The hub's list is the authority: a full sync replaces this one. + +It is an exact match on the content id. A file changed by one byte is another +id, and files over the partial-hash threshold are identified by a sample of their +bytes (§6.3) — a moderation tool, not a guarantee. +""" + +import json +import logging +import os +import re +from pathlib import Path + +log = logging.getLogger(__name__) + +_HASH = re.compile(r"^[0-9a-f]{64}$") + + +class ContentBlocklist: + def __init__(self, path: Path | None = None): + self._path = path + self._hashes: set[str] = set() + if path is not None: + try: + data = json.loads(path.read_text(encoding="utf-8")) + self._hashes = {h for h in data.get("hashes", []) if _HASH.match(h)} + except FileNotFoundError: + pass + except (OSError, ValueError) as e: + # A damaged file is not a reason to refuse to start; the next + # sync with the hub rewrites it. + log.warning("Content blocklist %s unreadable: %s", path, e) + + def __contains__(self, content_hash: object) -> bool: + return content_hash in self._hashes + + def __len__(self) -> int: + return len(self._hashes) + + def __iter__(self): + return iter(self._hashes) + + def replace(self, hashes) -> bool: + """The hub's full list. True when it differs from what was applied.""" + new = {h for h in hashes if isinstance(h, str) and _HASH.match(h)} + if new == self._hashes: + return False + self._hashes = new + self._save() + return True + + def apply(self, add=(), remove=()) -> bool: + """One pushed change. True when it changed anything.""" + before = set(self._hashes) + self._hashes |= {h for h in add if isinstance(h, str) and _HASH.match(h)} + self._hashes -= {h for h in remove if isinstance(h, str)} + if self._hashes == before: + return False + self._save() + return True + + def _save(self) -> None: + if self._path is None: + return + tmp = self._path.with_suffix(".tmp") + try: + tmp.write_text(json.dumps({"hashes": sorted(self._hashes)}), encoding="utf-8") + os.replace(tmp, self._path) + except OSError as e: + log.warning("Content blocklist %s not saved: %s", self._path, e) diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index a29aa9f..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") @@ -1448,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: @@ -1462,6 +1474,54 @@ class NodeDaemon(EnrichmentMixin): log.info("Index %s pushed to %d WebRTC peers", "delta" if delta is not None else "sync", pushed) + # ── 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).""" if not self._webrtc or not group_id: diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py index 5f47090..2ce8993 100644 --- a/packages/meshbay-node/src/meshbay_node/hub_client.py +++ b/packages/meshbay-node/src/meshbay_node/hub_client.py @@ -334,6 +334,8 @@ class HubClient: on_revocation: Any = None, on_webrtc_offer: Any = None, group_ids: list[str] | None = None, # static list or callable returning one + on_connected: Any = None, + on_blocklist: Any = None, ) -> None: """ Maintain a persistent WebSocket connection to the hub. @@ -393,6 +395,12 @@ class HubClient: self._ws = ws log.info("Hub WS connected") + if on_connected: + # Off the read loop, for the reason offers are: a slow + # fetch must not stop this socket being read. + task = asyncio.create_task(on_connected()) + pending.add(task) + task.add_done_callback(pending.discard) async for raw in ws: msg = json.loads(raw) @@ -405,6 +413,9 @@ class HubClient: elif mtype == "revocation" and on_revocation: on_revocation(msg.get("token", "")) + elif mtype == "blocklist_update" and on_blocklist: + on_blocklist(msg.get("add") or [], msg.get("remove") or []) + elif mtype == "webrtc_offer" and on_webrtc_offer: # Answered off the read loop on purpose. Awaiting the # handler here meant one slow negotiation stopped the @@ -460,6 +471,26 @@ class HubClient: log.warning("Could not deliver WebRTC answer to %s: %s", str(msg.get("peer_id"))[:8], e) + # ── Content blocklist (public groups) ───────────────────────────────── + + async def fetch_blocklist(self, max_pages: int = 100) -> set[str]: + """The hub's content blocklist, every page of it (docs/MESHBAY_DESIGN.md §7.5).""" + if self._session is None: + raise RuntimeError("Not logged in") + await self.ensure_fresh_token() + hashes: set[str] = set() + after = "" + for _ in range(max_pages): + r = await self._http.get("/v1/blocklist", params={"after": after}, + headers=self._session.auth_headers) + r.raise_for_status() + page = r.json() + hashes.update(page.get("hashes") or []) + after = page.get("next") or "" + if not after: + return hashes + raise RuntimeError("Content blocklist longer than expected; not applied") + # ── Convenience: full startup sequence ─────────────────────────────────── async def startup(self, endpoint_hint: str | None = None) -> HubSession: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py index f2b6fa5..857db33 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py @@ -64,6 +64,8 @@ class MusicMixin: """ ctx = self._group_ctx() file_id = msg.get("file_id", "") + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: self._send({"type": "error", "detail": "File not found"}) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py index a425c65..4337e24 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py @@ -228,6 +228,8 @@ class StreamingMixin: async def _stream_video_inner(self, msg: dict) -> None: ctx = self._group_ctx() file_id = msg.get("file_id", "") + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: self._send({"type": "error", "detail": "File not found"}) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py index bd2ab52..70781eb 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py @@ -39,6 +39,8 @@ class SubtitlesMixin: """ ctx = self._group_ctx() file_id = msg.get("file_id", "") + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: self._send({"type": "error", "detail": "File not found"}) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py index c3fb623..882ddbc 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py @@ -246,6 +246,25 @@ class SessionCore: return self._ctx["groups"].get(self._group_id) or {} return self._ctx + def _hidden_ids(self): + """The content blocklist, in a public group; nothing anywhere else. + + A private group's content never reaches the hub, so nothing in it can have + been blocked there (docs/MESHBAY_DESIGN.md §7.5, `meshbay_node.blocklist`). + """ + if self._group_ctx().get("visibility") != "public": + return () + return self._ctx.get("blocklist") or () + + def _refuse_blocked(self, file_id) -> bool: + """Refuse a file the blocklist names, in a public group. True if refused.""" + if not isinstance(file_id, str) or file_id not in self._hidden_ids(): + return False + self._send({"type": "error", "detail": "This file is not available here.", + "code": "content_blocked", "file_id": file_id}) + self._audit("content_blocked", file_id[:16]) + return True + def _register_peer(self) -> None: """Add this connection to its group's peer set. diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py index e76e23e..f629563 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py @@ -179,7 +179,7 @@ class FilesMixin: def _do_index_sync(self) -> None: ctx = self._group_ctx() - self._send(index_sync_message(ctx["index"], ctx.get("roots"))) + self._send(index_sync_message(ctx["index"], ctx.get("roots"), self._hidden_ids())) async def _try_serve_thumbnail( self, thumb_hash: str, chunk_index: int, gek: bytes | None, @@ -236,8 +236,18 @@ class FilesMixin: self._note_unleased(tr) file_id = msg["file_id"] chunk_index = msg["chunk_index"] + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: + # A blocked file's own thumbnail is a preview of it. + hidden = self._hidden_ids() + if hidden and file_id in { + e.thumb_hash for e in map(ctx["index"].get_entry, hidden) + if e is not None and e.thumb_hash}: + self._send({"type": "error", "detail": "This file is not available here.", + "code": "content_blocked", "file_id": file_id}) + return thumb = await self._try_serve_thumbnail(file_id, chunk_index, ctx.get("gek")) if thumb is not None: log.debug("file_req file_id=%s chunk=%s: served as thumbnail", diff --git a/packages/meshbay-node/src/meshbay_node/transport/wire.py b/packages/meshbay-node/src/meshbay_node/transport/wire.py index 1f9c925..258edb5 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/wire.py +++ b/packages/meshbay-node/src/meshbay_node/transport/wire.py @@ -65,7 +65,7 @@ def list_dirs(roots: RootSet | None) -> list[str]: return sorted(out)[:MAX_DIRS] -def index_sync_message(index, roots: RootSet | None) -> dict: +def index_sync_message(index, roots: RootSet | None, hidden=()) -> dict: """ The full `index_sync` message for one group. @@ -73,10 +73,13 @@ def index_sync_message(index, roots: RootSet | None) -> dict: without them a folder someone just created, or one they emptied, does not exist as far as a client is concerned, and a member cannot tell "the drive is unplugged" from "it is all still there". + + `hidden` is the content blocklist in a public group (`meshbay_node.blocklist`): + those entries are not sent at all. """ payload = { "version": index.version, - "entries": [index_entry_wire(e) for e in index.entries], + "entries": [index_entry_wire(e) for e in index.entries if e.id not in hidden], "dirs": list_dirs(roots), "roots": roots.describe() if roots else [], } @@ -88,7 +91,7 @@ def index_sync_message(index, roots: RootSet | None) -> dict: } -def index_delta_message(index, delta, roots=None) -> dict: +def index_delta_message(index, delta, roots=None, hidden=()) -> dict: """ One `index_delta` — what changed since the last thing this node broadcast. @@ -103,13 +106,16 @@ def index_delta_message(index, delta, roots=None) -> dict: page: the delta that told them something had changed was the one message that could not say what. It is a handful of dicts, bounded by the number of directories a group has, and it is sealed with the rest. + + `hidden`, as for `index_sync_message`: a blocked entry is never added or + updated; its deletion still goes out. """ payload = { "base_version": delta.base_version, "version": delta.version, - "additions": [index_entry_wire(e) for e in delta.additions], + "additions": [index_entry_wire(e) for e in delta.additions if e.id not in hidden], "deletions": list(delta.deletions), - "updates": [index_entry_wire(e) for e in delta.updates], + "updates": [index_entry_wire(e) for e in delta.updates if e.id not in hidden], } if roots is not None: payload["roots"] = roots.describe() -- cgit v1.2.3