""" 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