1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
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
|