diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-09-28 21:49:14 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-09-28 21:49:14 +0200 |
| commit | 07480eb3f8ad0bb4369ac8c41df7c4140b108d0e (patch) | |
| tree | 2f692c0dccb184890fd82339590a4e9ce168c6e4 /packages | |
| parent | 91505face56f7ee6817408e52bad7902add75f09 (diff) | |
| download | meshbay-07480eb3f8ad0bb4369ac8c41df7c4140b108d0e.tar.gz | |
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 <noreply@anthropic.com>
Diffstat (limited to 'packages')
15 files changed, 627 insertions, 55 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/deps.py b/packages/meshbay-hub/src/meshbay_hub/api/deps.py index 42f4101..2fd8b99 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/deps.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/deps.py @@ -90,6 +90,17 @@ async def require_user_scope( return current_user +async def require_node_scope( + payload: dict = Depends(_decode_token), + current_user: User = Depends(get_current_user), +) -> User: + """Only a node daemon's own token — for what a node fetches on its own behalf.""" + if payload.get("scope") != "node": + raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, + detail="Node token required") + return current_user + + def user_is_admin(user: User) -> bool: """Admin by DB role or by the config allow-list. Use inside a handler that already depends on `require_moderator` but has to draw the admin line for diff --git a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py index 35c688c..c4e6af5 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py @@ -18,8 +18,9 @@ Admin endpoints: DELETE /v1/admin/blocklist/{hash} — remove a hash Node integration: - GET /v1/blocklist/check?hash=<blake3> — check if a hash is blocked - GET /v1/blocklist — full blocklist (for node sync) + GET /v1/blocklist?after=<hash> — the list, paged, for a node's own token + WebSocket `blocklist_update` — additions and removals, pushed to nodes + hosting a public group """ import logging @@ -30,9 +31,10 @@ from sqlalchemy import func, select from sqlalchemy.ext.asyncio import AsyncSession from meshbay_hub import hub_settings -from meshbay_hub.api.deps import get_current_user, require_admin +from meshbay_hub.api.deps import get_current_user, require_admin, require_node_scope from meshbay_hub.api.middleware import limiter from meshbay_hub.api.netutil import client_ip +from meshbay_hub.api.revocation import broadcast_blocklist_update from meshbay_hub.db.engine import get_db from meshbay_hub.db.models import ContentBlocklist, ContentReport, User @@ -47,6 +49,10 @@ router = APIRouter(tags=["moderation"]) AUTO_BLOCK_THRESHOLD = 3 +def _is_hash(value: str) -> bool: + return len(value) == 64 and all(c in "0123456789abcdef" for c in value) + + # ── Models ──────────────────────────────────────────────────────────────────── class ReportRequest(BaseModel): @@ -89,7 +95,7 @@ async def report_content( status_code=403, detail="This hub does not broker public content, so there is nothing to report here.") - if len(body.content_hash) != 64 or not all(c in "0123456789abcdef" for c in body.content_hash): + if not _is_hash(body.content_hash): raise HTTPException(status_code=422, detail="content_hash must be 64 hex chars (blake3)") # One vote per account per hash — a single reporter must not be able to walk @@ -115,6 +121,7 @@ async def report_content( .where(ContentReport.content_hash == body.content_hash)) or 0 action = "already_reported" if already else "logged" + auto_blocked = False if distinct_reporters >= AUTO_BLOCK_THRESHOLD: existing = await db.get(ContentBlocklist, body.content_hash) if not existing: @@ -124,10 +131,13 @@ async def report_content( added_by="auto", )) action = "auto_blocked" + auto_blocked = True log.warning("Content auto-blocked after %d distinct reporters: %s", distinct_reporters, body.content_hash[:16]) await db.commit() + if auto_blocked: + await broadcast_blocklist_update(db, add=[body.content_hash]) return { "status": action, "content_hash": body.content_hash, @@ -136,49 +146,31 @@ async def report_content( } -@router.get("/v1/blocklist/check") -@limiter.limit("120/minute") -async def check_blocklist( - hash: str, - request: Request, - db: AsyncSession = Depends(get_db), -): - """Check if a single hash is blocked. Used by nodes before serving public content. - - Unauthenticated, because a node consults it before serving public content - and does so on its own behalf. That makes the shape check worth having: - without it any string of any length became a primary-key lookup. - """ - if len(hash) != 64 or not all(c in "0123456789abcdef" for c in hash): - raise HTTPException(status_code=422, detail="hash must be 64 hex chars (blake3)") - blocked = await db.get(ContentBlocklist, hash) - return { - "blocked": blocked is not None, - "hash": hash, - "reason": blocked.reason if blocked else None, - } - - @router.get("/v1/blocklist") async def get_blocklist( + current_node: User = Depends(require_node_scope), db: AsyncSession = Depends(get_db), - # Bounded, like every other list. This one takes no authentication — a - # node syncs it at startup — and had no ceiling at all, so any stranger - # could ask for the table in one query, repeatedly. 10 000 is what a node - # asks for, so it is the default and also the most anyone may have. + after: str = Query(default="", max_length=64), limit: int = Query(default=10000, ge=1, le=10000), ): """ - Return the full blocklist. Nodes sync this on startup. - Returns hashes only (not reasons) to minimize data exposure. + The content blocklist, a page at a time, for a node hosting a public group. + + Hashes only, never the reasons. Ordered by hash so `after` (the last hash of + the previous page) is a stable cursor: a list longer than one page used to + be cut at 10 000 with no way to ask for the rest, and the node applying it + silently served everything past the cut. A node's own token, because this + is what a node fetches on its own behalf and nothing else asks for it. """ result = await db.execute( select(ContentBlocklist.content_hash) - .order_by(ContentBlocklist.added_at.desc()) + .where(ContentBlocklist.content_hash > after) + .order_by(ContentBlocklist.content_hash) .limit(limit) ) hashes = [row[0] for row in result.fetchall()] - return {"count": len(hashes), "hashes": hashes} + return {"hashes": hashes, + "next": hashes[-1] if len(hashes) == limit else None} # ── Admin endpoints ─────────────────────────────────────────────────────────── @@ -214,16 +206,19 @@ async def admin_add_blocklist( current_user: User = Depends(require_admin), db: AsyncSession = Depends(get_db), ): + if not _is_hash(body.content_hash): + raise HTTPException(status_code=422, detail="content_hash must be 64 hex chars (blake3)") existing = await db.get(ContentBlocklist, body.content_hash) if existing: raise HTTPException(status_code=409, detail="Hash already blocked") db.add(ContentBlocklist( content_hash=body.content_hash, - reason=body.reason, + reason=body.reason[:64], added_by=current_user.username, )) await db.commit() + await broadcast_blocklist_update(db, add=[body.content_hash]) return {"status": "blocked", "hash": body.content_hash} @@ -238,5 +233,6 @@ async def admin_remove_blocklist( raise HTTPException(status_code=404, detail="Hash not in blocklist") await db.delete(entry) await db.commit() + await broadcast_blocklist_update(db, remove=[content_hash]) return {"status": "unblocked", "hash": content_hash} diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py index 2c0b8db..fe9edf3 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py @@ -165,6 +165,35 @@ async def broadcast_revocation(token: str) -> int: return sent +async def broadcast_blocklist_update(db: AsyncSession, *, add: list[str] = (), + remove: list[str] = ()) -> int: + """ + Tell every connected node that hosts a public group what changed on the + content blocklist. Returns how many were told. + + Only those nodes: the list names public content, and a node hosting only + private groups has nothing to apply it to. A node that is offline now syncs + the whole list when it next connects (`GET /v1/blocklist`), so a missed push + is only late, never lost. + """ + if not add and not remove: + return 0 + public = set((await db.execute( + select(Group.id).where(Group.visibility == "public"))).scalars().all()) + payload = json.dumps({"type": "blocklist_update", + "add": list(add), "remove": list(remove)}) + sent = 0 + for node_id, ws in list(_connected_nodes.items()): + if not public.intersection(_node_groups.get(node_id, [])): + continue + try: + await ws.send_text(payload) + sent += 1 + except Exception: + _connected_nodes.pop(node_id, None) + return sent + + def _sign_revocation(target: str, target_id: str, reason: str) -> str: """Issue a signed revocation token (JWT EdDSA).""" from meshbay_hub.auth import _hub_id, _hub_sk_pem diff --git a/packages/meshbay-hub/tests/test_moderation.py b/packages/meshbay-hub/tests/test_moderation.py index 3109e29..b328d67 100644 --- a/packages/meshbay-hub/tests/test_moderation.py +++ b/packages/meshbay-hub/tests/test_moderation.py @@ -22,6 +22,20 @@ async def _register_and_login(client, username: str) -> dict: return {"Authorization": f"Bearer {r.json()['access_token']}"} +async def _node_headers(client, username: str) -> dict: + """A node daemon's token for a fresh account — what a node syncs with.""" + from meshbay_hub.auth import issue_access_token + user = await _register_and_login(client, username) + me = (await client.get("/v1/users/me", headers=user)).json() + tok = issue_access_token(me["user_id"], ttl=3600, groups=[], scope="node") + return {"Authorization": f"Bearer {tok}"} + + +async def _blocked(client, h: str) -> bool: + node = await _node_headers(client, f"node_{h[:6]}_{len(h)}") + return h in (await client.get("/v1/blocklist", headers=node)).json()["hashes"] + + @pytest.fixture async def reporter(client): return await _register_and_login(client, "reporter_one") @@ -69,8 +83,7 @@ async def test_same_reporter_cannot_walk_the_threshold(client, reporter): assert r.json()["report_count"] == 1 assert r.json()["status"] == "already_reported" - check = await client.get(f"/v1/blocklist/check?hash={h}") - assert check.json()["blocked"] is False + assert not await _blocked(client, h) @pytest.mark.asyncio @@ -84,8 +97,7 @@ async def test_auto_block_on_distinct_reporters(client): assert r.json()["status"] == "auto_blocked" assert r.json()["report_count"] == 3 - check = await client.get(f"/v1/blocklist/check?hash={h}") - assert check.json()["blocked"] is True + assert await _blocked(client, h) @pytest.mark.asyncio @@ -118,14 +130,13 @@ async def test_admin_add_remove_blocklist(client, admin_headers): headers=admin_headers) assert r.status_code == 201 - r = await client.get(f"/v1/blocklist/check?hash={hash4}") - assert r.json()["blocked"] is True + assert await _blocked(client, hash4) r = await client.delete(f"/v1/admin/blocklist/{hash4}", headers=admin_headers) assert r.status_code == 200 - r = await client.get(f"/v1/blocklist/check?hash={hash4}") - assert r.json()["blocked"] is False + node = await _node_headers(client, "node_after_unblock") + assert hash4 not in (await client.get("/v1/blocklist", headers=node)).json()["hashes"] @pytest.mark.asyncio @@ -134,6 +145,80 @@ async def test_full_blocklist(client, admin_headers): await client.post("/v1/admin/blocklist", json={"content_hash": hash5, "reason": "test"}, headers=admin_headers) - r = await client.get("/v1/blocklist") + node = await _node_headers(client, "node_full") + r = await client.get("/v1/blocklist", headers=node) assert r.status_code == 200 assert hash5 in r.json()["hashes"] + + +@pytest.mark.asyncio +async def test_only_a_node_token_reads_the_blocklist(client, reporter): + assert (await client.get("/v1/blocklist")).status_code in (401, 422) + assert (await client.get("/v1/blocklist", headers=reporter)).status_code == 403 + + +@pytest.mark.asyncio +async def test_the_blocklist_pages_past_its_limit(client, admin_headers): + hashes = sorted(f"{i:064x}" for i in range(5)) + for h in hashes: + await client.post("/v1/admin/blocklist", + json={"content_hash": h, "reason": "test"}, + headers=admin_headers) + node = await _node_headers(client, "node_pager") + seen, after = [], "" + while True: + page = (await client.get(f"/v1/blocklist?limit=2&after={after}", + headers=node)).json() + seen += page["hashes"] + if not page["next"]: + break + after = page["next"] + assert [h for h in seen if h in hashes] == hashes + + +@pytest.mark.asyncio +async def test_an_admin_cannot_block_something_that_is_not_a_hash(client, admin_headers): + r = await client.post("/v1/admin/blocklist", + json={"content_hash": "x" * 5000, "reason": "test"}, + headers=admin_headers) + assert r.status_code == 422 + + +@pytest.mark.asyncio +async def test_a_change_is_pushed_to_nodes_hosting_a_public_group_only(db_session): + """A node hosting only private groups has nothing to apply the list to.""" + import json + + import meshbay_hub.api.revocation as rev + from meshbay_hub.db.models import Group, User + + db_session.add(User(id="owner-bl", username="owner_bl", email="x", hub_id="h", + pw_hash=b"x", pw_salt=b"x")) + db_session.add(Group(id="pub-bl", name="pub", admin_id="owner-bl", + visibility="public")) + db_session.add(Group(id="priv-bl", name="priv", admin_id="owner-bl", + visibility="private")) + await db_session.commit() + + class _WS: + def __init__(self): + self.sent = [] + + async def send_text(self, text): + self.sent.append(json.loads(text)) + + public_node, private_node = _WS(), _WS() + saved = dict(rev._connected_nodes), dict(rev._node_groups) + rev._connected_nodes.update({"n-pub": public_node, "n-priv": private_node}) + rev._node_groups.update({"n-pub": ["pub-bl"], "n-priv": ["priv-bl"]}) + try: + sent = await rev.broadcast_blocklist_update(db_session, add=["a" * 64]) + finally: + rev._connected_nodes.clear() + rev._connected_nodes.update(saved[0]) + rev._node_groups.clear() + rev._node_groups.update(saved[1]) + assert sent == 1 + assert public_node.sent == [{"type": "blocklist_update", "add": ["a" * 64], + "remove": []}] + assert private_node.sent == [] diff --git a/packages/meshbay-hub/tests/test_unauthenticated_surface.py b/packages/meshbay-hub/tests/test_unauthenticated_surface.py index f2da42b..78615ca 100644 --- a/packages/meshbay-hub/tests/test_unauthenticated_surface.py +++ b/packages/meshbay-hub/tests/test_unauthenticated_surface.py @@ -42,8 +42,6 @@ PUBLIC = { ("POST", "/v1/nodes/auth"): "node sign-in — Ed25519 signature over a fresh timestamp", ("WS", "/v1/nodes/ws"): "a node-scoped JWT in the first message, within a timeout", ("GET", "/v1/groups"): "the public directory — empty when public groups are off", - ("GET", "/v1/blocklist"): "nodes sync it on their own behalf — hashes only", - ("GET", "/v1/blocklist/check"): "nodes consult it on their own behalf", ("GET", "/v1/relays"): "closed: 503 while relay.RELAYS_ENABLED is False", ("POST", "/v1/relays/register"): "closed; when open, approved key + signature", ("GET", "/mhp/info"): "closed: 503 while federation.FEDERATION_ENABLED is False", 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() diff --git a/packages/meshbay-node/tests/test_content_blocklist.py b/packages/meshbay-node/tests/test_content_blocklist.py new file mode 100644 index 0000000..0e07787 --- /dev/null +++ b/packages/meshbay-node/tests/test_content_blocklist.py @@ -0,0 +1,238 @@ +""" +The hub's content blocklist, applied by a node in its public groups +(docs/MESHBAY_DESIGN.md §7.5, `meshbay_node.blocklist`). + +A blocked file leaves the index members are sent and is refused if asked for — +in a public group, and nowhere else: a private group's content never reaches the +hub, so nothing there can have been blocked. +""" + +import struct + +import msgpack +import pytest +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey +from meshbay_common.crypto import generate_gek +from meshbay_common.groupbox import PURPOSE_INDEX, unseal +from meshbay_common.protocol import IndexEntry +from meshbay_node.blocklist import ContentBlocklist +from meshbay_node.hub_client import HubClient +from meshbay_node.indexer import GroupIndex +from meshbay_node.transport.webrtc_server import WebRTCPeerSession +from meshbay_node.transport.wire import index_delta_message, index_sync_message + +GROUP = "g" * 32 +BLOCKED = "b" * 64 +KEPT = "c" * 64 +THUMB = "d" * 64 + + +# ── The list ────────────────────────────────────────────────────────────────── + +def test_the_list_survives_a_restart(tmp_path): + path = tmp_path / "blocklist.json" + assert ContentBlocklist(path).replace([BLOCKED]) + assert BLOCKED in ContentBlocklist(path) + + +def test_a_full_sync_replaces_and_says_whether_anything_changed(tmp_path): + bl = ContentBlocklist(tmp_path / "blocklist.json") + assert bl.replace([BLOCKED, KEPT]) + assert not bl.replace([KEPT, BLOCKED]) + assert bl.replace([KEPT]) + assert BLOCKED not in bl and KEPT in bl + + +def test_a_pushed_change_applies_and_ignores_what_is_not_a_hash(tmp_path): + bl = ContentBlocklist(tmp_path / "blocklist.json") + assert bl.apply(add=[BLOCKED, "../etc", 7, "B" * 64]) + assert list(bl) == [BLOCKED] + assert not bl.apply(add=[BLOCKED]) + assert bl.apply(remove=[BLOCKED]) + assert len(bl) == 0 + + +def test_a_damaged_file_does_not_stop_the_node(tmp_path): + path = tmp_path / "blocklist.json" + path.write_text("{not json", encoding="utf-8") + assert len(ContentBlocklist(path)) == 0 + + +# ── What members are sent ───────────────────────────────────────────────────── + +def _index(gek) -> GroupIndex: + idx = GroupIndex(group_id=GROUP, sk_node=Ed25519PrivateKey.generate(), gek=gek) + idx.add_entry(IndexEntry(id=BLOCKED, name="a.jpg", path="r", size=1, type="image", + added_at=0, thumb_hash=THUMB)) + idx.add_entry(IndexEntry(id=KEPT, name="b.jpg", path="r", size=1, type="image", + added_at=0)) + return idx + + +def test_a_hidden_entry_is_not_in_the_index_or_its_deltas(): + gek = generate_gek() + idx = _index(gek) + msg = index_sync_message(idx, None, {BLOCKED}) + ids = [e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]] + assert ids == [KEPT] + + delta = idx.diff(GroupIndex(group_id=GROUP, sk_node=idx.sk_node, gek=gek, version=0)) + payload = unseal(gek, PURPOSE_INDEX, "index_delta", GROUP, + index_delta_message(idx, delta, None, {BLOCKED})) + assert [e["id"] for e in payload["additions"]] == [KEPT] + + +# ── What a session serves ──────────────────────────────────────────────────── + +class _Channel: + readyState = "open" + + def __init__(self): + self.sent = [] + + def send(self, data: bytes) -> None: + (n,) = struct.unpack(">I", data[:4]) + self.sent.append(msgpack.unpackb(data[4:4 + n], raw=False)) + + +class _PC: + connectionState = "connected" + iceConnectionState = "connected" + remoteDescription = None + localDescription = None + sctp = None + + +def _session(visibility: str, tmp_path): + gek = generate_gek() + bl = ContentBlocklist(tmp_path / "blocklist.json") + bl.replace([BLOCKED]) + ctx = {"sk_node": Ed25519PrivateKey.from_private_bytes(b"\x01" * 32), + "groups": {GROUP: {"visibility": visibility, "index": _index(gek), + "gek": gek, "roots": None}}, + "blocklist": bl} + s = WebRTCPeerSession(_PC(), ctx, peer_id="peer") + s._channel = _Channel() + s._audit = lambda *a, **k: None + s._user_id, s._group_id = "member-1", GROUP + return s, gek + + +@pytest.mark.asyncio +async def test_a_public_group_refuses_a_blocked_file_and_its_thumbnail(tmp_path): + s, _ = _session("public", tmp_path) + for file_id in (BLOCKED, THUMB): + await s._do_file_request({"file_id": file_id, "chunk_index": 0}) + assert s._channel.sent[-1]["code"] == "content_blocked" + + +@pytest.mark.asyncio +async def test_a_public_group_does_not_list_a_blocked_file(tmp_path): + s, gek = _session("public", tmp_path) + s._do_index_sync() + payload = unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, s._channel.sent[-1]) + assert [e["id"] for e in payload["entries"]] == [KEPT] + + +@pytest.mark.asyncio +async def test_a_private_group_is_untouched(tmp_path): + s, gek = _session("private", tmp_path) + s._do_index_sync() + payload = unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, s._channel.sent[-1]) + assert sorted(e["id"] for e in payload["entries"]) == [BLOCKED, KEPT] + assert not s._refuse_blocked(BLOCKED) + + +@pytest.mark.asyncio +async def test_streaming_subtitles_and_transcoding_are_refused_too(tmp_path): + s, _ = _session("public", tmp_path) + await s._stream_video_inner({"file_id": BLOCKED}) + assert s._channel.sent[-1]["code"] == "content_blocked" + await s._do_subtitle_request({"file_id": BLOCKED, "track": 0}) + assert s._channel.sent[-1]["code"] == "content_blocked" + await s._do_audio_transcode_request({"file_id": BLOCKED}) + assert s._channel.sent[-1]["code"] == "content_blocked" + + +# ── How the node gets the list ─────────────────────────────────────────────── + +@pytest.mark.asyncio +async def test_the_whole_list_is_fetched_page_by_page(): + pages = {"": {"hashes": [BLOCKED], "next": BLOCKED}, + BLOCKED: {"hashes": [KEPT], "next": None}} + + class _Resp: + def __init__(self, body): + self._body = body + + def raise_for_status(self): + pass + + def json(self): + return self._body + + class _Http: + async def get(self, path, params, headers): + assert path == "/v1/blocklist" + return _Resp(pages[params["after"]]) + + class _Session: + auth_headers = {} + + hub = HubClient.__new__(HubClient) + hub._session, hub._http = _Session(), _Http() + + async def fresh(): + return None + hub.ensure_fresh_token = fresh + assert await hub.fetch_blocklist() == {BLOCKED, KEPT} + + +# ── A change reaches members already connected ─────────────────────────────── + +def _daemon(tmp_path, visibility: str): + from unittest.mock import MagicMock + + from meshbay_node.config import Config, GroupConfig, HubConfig, KeystoreConfig, NodeConfig + from meshbay_node.daemon import NodeDaemon + + shared = tmp_path / "shared" + shared.mkdir() + config = Config( + hub=HubConfig(url="http://localhost:9999", username="t"), + node=NodeConfig(), + groups=[GroupConfig(id=GROUP, name="g", shared_dir=str(shared), + visibility=visibility)], + keystore=KeystoreConfig(path=tmp_path / "keystore.enc"), + data_dir=tmp_path / "data", + ) + daemon = NodeDaemon(config) + gek = generate_gek() + indexer = MagicMock() + indexer.index, indexer.roots = _index(gek), None + daemon._indexers = [indexer] + session = MagicMock() + session._group_id = GROUP + daemon._webrtc = MagicMock() + daemon._webrtc._sessions = {"p": session} + return daemon, session, gek + + +def test_a_pushed_block_resends_the_index_without_the_file(tmp_path): + daemon, session, gek = _daemon(tmp_path, "public") + daemon._on_blocklist_update([BLOCKED], []) + msg = session._send.call_args[0][0] + ids = [e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]] + assert ids == [KEPT] + + daemon._on_blocklist_update([], [BLOCKED]) + msg = session._send.call_args[0][0] + ids = sorted(e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]) + assert ids == [BLOCKED, KEPT] + + +def test_a_node_with_only_private_groups_ignores_the_list(tmp_path): + daemon, session, _ = _daemon(tmp_path, "private") + daemon._on_blocklist_update([BLOCKED], []) + session._send.assert_not_called() + assert BLOCKED not in daemon._blocklist |