diff options
Diffstat (limited to 'packages/meshbay-node/tests')
| -rw-r--r-- | packages/meshbay-node/tests/golden/dispatch.json | 448 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_admin_challenge_bounds.py | 135 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_content_blocklist.py | 238 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_daemon.py | 74 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_indexer.py | 96 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_security_regressions.py | 21 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_tmdb_config_policy.py | 13 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_webrtc_transport.py | 10 | ||||
| -rwxr-xr-x | packages/meshbay-node/tests/transfer_probe.py | 29 |
9 files changed, 488 insertions, 576 deletions
diff --git a/packages/meshbay-node/tests/golden/dispatch.json b/packages/meshbay-node/tests/golden/dispatch.json index 5d15057..a7bb543 100644 --- a/packages/meshbay-node/tests/golden/dispatch.json +++ b/packages/meshbay-node/tests/golden/dispatch.json @@ -1874,166 +1874,6 @@ "sent": [], "spawned": [] }, - "chat_attach | challenged | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | challenged | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | challenged | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | challenged | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | member | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | member | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | member | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | member | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, "chat_directory | challenged | bare": { "audit": [], "log": [], @@ -2228,7 +2068,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2551,7 +2391,7 @@ "subject": "gggggggggggggggggggggggggggggggg", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2570,7 +2410,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2589,7 +2429,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2608,7 +2448,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2878,7 +2718,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2892,7 +2732,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2906,7 +2746,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2920,7 +2760,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2934,7 +2774,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2948,7 +2788,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2962,7 +2802,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2976,7 +2816,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -7229,166 +7069,6 @@ "sent": [], "spawned": [] }, - "ephemeral_stream | challenged | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | challenged | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | challenged | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | challenged | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | member | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | member | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | member | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | member | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, "file_chunk | challenged | bare": { "audit": [], "log": [], @@ -8855,7 +8535,7 @@ "subject": "gggggggggggggggggggggggggggggggg", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -8874,7 +8554,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -8893,7 +8573,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -8912,7 +8592,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9244,10 +8924,10 @@ "op": "group_attach", "op_id": "<volatile>", "req_id": 4242, - "subject": "['x']", + "subject": "{\"name\":\"['x']\",\"shared_dir\":\"['x']\",\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9263,10 +8943,10 @@ "op": "group_attach", "op_id": "<volatile>", "req_id": 4242, - "subject": "7", + "subject": "{\"name\":\"7\",\"shared_dir\":\"7\",\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9282,10 +8962,10 @@ "op": "group_attach", "op_id": "<volatile>", "req_id": 4242, - "subject": "x", + "subject": "{\"name\":\"x\",\"shared_dir\":\"x\",\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9620,7 +9300,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9639,7 +9319,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9658,7 +9338,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -11961,7 +11641,7 @@ "subject": "link:gggggggggggggggggggggggggggggggg", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14060,7 +13740,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14079,7 +13759,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14098,7 +13778,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14433,7 +14113,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14452,7 +14132,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14471,7 +14151,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16532,7 +16212,7 @@ "req_id": 4242, "token": null, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16547,7 +16227,7 @@ "x" ], "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16560,7 +16240,7 @@ "req_id": 4242, "token": 7, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16573,7 +16253,7 @@ "req_id": 4242, "token": "x", "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16586,7 +16266,7 @@ "req_id": 4242, "token": null, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16601,7 +16281,7 @@ "x" ], "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16614,7 +16294,7 @@ "req_id": 4242, "token": 7, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16627,7 +16307,7 @@ "req_id": 4242, "token": "x", "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16959,10 +16639,10 @@ "op": "root_add", "op_id": "<volatile>", "req_id": 4242, - "subject": "['x']", + "subject": "{\"kind\":\"['x']\",\"name\":\"['x']\",\"path\":\"['x']\",\"removable\":true,\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16978,10 +16658,10 @@ "op": "root_add", "op_id": "<volatile>", "req_id": 4242, - "subject": "7", + "subject": "{\"kind\":\"7\",\"name\":\"7\",\"path\":\"7\",\"removable\":true,\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16997,10 +16677,10 @@ "op": "root_add", "op_id": "<volatile>", "req_id": 4242, - "subject": "x", + "subject": "{\"kind\":\"x\",\"name\":\"x\",\"path\":\"x\",\"removable\":true,\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17335,7 +17015,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17354,7 +17034,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17373,7 +17053,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17708,7 +17388,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17727,7 +17407,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17746,7 +17426,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18081,7 +17761,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18100,7 +17780,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18119,7 +17799,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18454,7 +18134,7 @@ "subject": "['x']:rw=on,rem=on", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18473,7 +18153,7 @@ "subject": "7:rw=on,rem=on", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18492,7 +18172,7 @@ "subject": "x:rw=on,rem=on", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -21436,10 +21116,10 @@ "op": "tmdb_config", "op_id": "<volatile>", "req_id": 4242, - "subject": "custom_token=no,language=default", + "subject": "{\"language\":null,\"token\":null}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -21479,10 +21159,10 @@ "op": "tmdb_config", "op_id": "<volatile>", "req_id": 4242, - "subject": "custom_token=yes,language=x", + "subject": "{\"language\":\"x\",\"token\":\"sha256:2d711642b726b04401627ca9fbac32f5c8530fb1903cc4db02258717921a4881\"}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -23369,7 +23049,7 @@ "subject": "d=7,u=7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] diff --git a/packages/meshbay-node/tests/test_admin_challenge_bounds.py b/packages/meshbay-node/tests/test_admin_challenge_bounds.py new file mode 100644 index 0000000..fd3b2f1 --- /dev/null +++ b/packages/meshbay-node/tests/test_admin_challenge_bounds.py @@ -0,0 +1,135 @@ +""" +What a connection may leave waiting for a signature (docs/MESHBAY_DESIGN.md §13.5b). + +Anyone authenticated can ask for an admin challenge — the signature is checked +later — so a member who never answers must not make the node keep every request. +Measured before the bound: 200 `root_add` of 1 MiB each from a plain member held +200 pending operations and ~400 MiB for the life of the connection. +""" + +import struct +import time + +import msgpack +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey +from meshbay_node.transport.webrtc.admin import MAX_ADMIN_OP_BYTES, MAX_PENDING_ADMIN_OPS +from meshbay_node.transport.webrtc_server import WebRTCPeerSession + +GROUP = "g" * 32 + + +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 _member_session(): + """An authenticated member — not the operator — on a node that has one.""" + ctx = {"sk_node": Ed25519PrivateKey.from_private_bytes(b"\x01" * 32), + "groups": {GROUP: {}}, "has_admin_authority": True} + 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 + + +def _root_add(s, path: str) -> dict: + s._dispatch_message({"type": "root_add", "group_id": GROUP, "path": path}) + return s._channel.sent[-1] + + +def test_a_member_cannot_pile_up_challenges(): + s = _member_session() + for i in range(MAX_PENDING_ADMIN_OPS): + assert _root_add(s, f"/srv/{i}")["type"] == "admin_challenge" + refused = _root_add(s, "/srv/one-too-many") + assert refused["type"] == "error" and refused["code"] == "too_many_pending" + assert len(s._admin_ops) == MAX_PENDING_ADMIN_OPS + + +def test_an_oversized_request_is_not_kept(): + s = _member_session() + refused = _root_add(s, "x" * (MAX_ADMIN_OP_BYTES + 1)) + assert refused["type"] == "error" and refused["code"] == "too_large" + assert s._admin_ops == {} + + +def test_an_expired_challenge_frees_its_place(): + s = _member_session() + for i in range(MAX_PENDING_ADMIN_OPS): + _root_add(s, f"/srv/{i}") + for pending in s._admin_ops.values(): + pending["ts"] -= 10_000 + assert _root_add(s, "/srv/after-expiry")["type"] == "admin_challenge" + assert len(s._admin_ops) == 1 + + +def test_answering_a_challenge_frees_its_place(): + s = _member_session() + for i in range(MAX_PENDING_ADMIN_OPS): + _root_add(s, f"/srv/{i}") + op_id = next(iter(s._admin_ops)) + s._dispatch_message({"type": "admin_response", "op_id": op_id, "signature": "!!"}) + assert len(s._admin_ops) == MAX_PENDING_ADMIN_OPS - 1 + assert _root_add(s, "/srv/next")["type"] == "admin_challenge" + assert all(time.time() - p["ts"] < 5 for p in s._admin_ops.values()) + + +# ── What a challenge covers (docs/MESHBAY_DESIGN.md §5.4) ──────────────────── +# +# The signature covers the subject and nothing else of a request, so every value +# the executor acts on has to be in it. + +def test_root_add_signs_whether_members_may_write(): + from meshbay_common.adminop import root_add_subject + s = _member_session() + s._dispatch_message({"type": "root_add", "group_id": GROUP, "path": "/srv/drop", + "name": "Drop", "writable": True, "removable": False}) + challenge = s._channel.sent[-1] + assert challenge["subject"] == root_add_subject("/srv/drop", "Drop", "generic", + True, False) + assert challenge["subject"] != root_add_subject("/srv/drop", "Drop", "generic", + False, False) + + +def test_group_attach_signs_the_directory_it_exposes(): + from meshbay_common.adminop import group_attach_subject + s = _member_session() + s._dispatch_message({"type": "group_attach", "name": "photos", + "shared_dir": "/home/me/Photos"}) + assert s._channel.sent[-1]["subject"] == group_attach_subject( + "photos", "/home/me/Photos", True) + + +def test_invite_create_signs_the_name_it_records(): + from meshbay_common.adminop import invite_create_subject + s = _member_session() + s._ctx["roster"] = object() # only its presence is checked before the challenge + s._dispatch_message({"type": "invite_create", "group_id": GROUP, + "user_id": "u-1", "username": "alice"}) + assert s._channel.sent[-1]["subject"] == invite_create_subject("u-1", "alice") + + +def test_tmdb_config_signs_the_token_without_writing_it(): + from meshbay_common.adminop import tmdb_config_subject + s = _member_session() + s._dispatch_message({"type": "tmdb_config", "token": "secret-token", + "language": "fr-FR"}) + subject = s._channel.sent[-1]["subject"] + assert subject == tmdb_config_subject("secret-token", "fr-FR") + assert "secret-token" not in subject 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 diff --git a/packages/meshbay-node/tests/test_daemon.py b/packages/meshbay-node/tests/test_daemon.py index 8c4da2d..acaafac 100644 --- a/packages/meshbay-node/tests/test_daemon.py +++ b/packages/meshbay-node/tests/test_daemon.py @@ -2,7 +2,7 @@ Integration test: Node daemon wires all components correctly. Phase 11 — verifies that NodeDaemon creates chat stores, WebRTC transport, -index push on change, swarm registration, and shuts down cleanly. +index push on change, and shuts down cleanly. Hub interaction is mocked. """ @@ -241,7 +241,6 @@ async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hu daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 # real value would make this test wait 0.5s daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=2) daemon._state["endpoint_hint"] = "node123" sk_node = Ed25519PrivateKey.generate() @@ -268,47 +267,10 @@ async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hu payload = unseal(gek, PURPOSE_INDEX, "index_sync", "a" * 32, msg) assert len(payload["entries"]) == indexer.index.count - # Finding H7: this group is private, so its content hashes must NOT be - # registered with the hub. The test previously asserted the opposite — - # publishing a fingerprint of every private file was treated as expected - # behaviour. Index push to members is unaffected (asserted above). + # Finding H7: a change to the index tells the hub nothing — no content hash + # of any group reaches it. Index push to members is unaffected (above). await asyncio.sleep(0.1) - daemon._hub.register_swarm.assert_not_called() - -@pytest.mark.asyncio -async def test_daemon_index_change_registers_swarm_for_public_group( - tmp_path, shared_dir, gek, hub_pk_pem): - """Public groups still register content hashes with the hub swarm (H7).""" - config = Config( - hub=HubConfig(url="http://localhost:9999", username="testuser"), - node=NodeConfig(quic_port=_free_port(), ui_port=_free_port()), - groups=[GroupConfig( - id="a" * 32, - name="public-group", - shared_dir=str(shared_dir), - visibility="public", - quic_port=29010, - )], - keystore=KeystoreConfig(path=tmp_path / "keystore.enc"), - data_dir=tmp_path / "data", - ) - daemon = NodeDaemon(config) - daemon._broadcast_coalesce_secs = 0.01 - daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=2) - daemon._state["endpoint_hint"] = "node123" - - indexer = DirectoryIndexer( - roots=one_root(shared_dir), group_id="a" * 32, - sk_node=Ed25519PrivateKey.generate(), gek=gek) - await indexer.initial_scan() - - await daemon._on_index_change(indexer) - - await asyncio.sleep(0.1) - daemon._hub.register_swarm.assert_called_once() - assert len(daemon._hub.register_swarm.call_args[0][0]) == indexer.index.count - + assert daemon._hub.mock_calls == [] @pytest.mark.asyncio async def test_daemon_index_change_skips_other_group_peers( @@ -325,7 +287,6 @@ async def test_daemon_index_change_skips_other_group_peers( daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=0) daemon._state["endpoint_hint"] = "node123" sk_node = Ed25519PrivateKey.generate() @@ -370,7 +331,6 @@ def _new_daemon_for_group(tmp_path, shared_dir, gek, group_id="a" * 32, daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=0) daemon._state["endpoint_hint"] = "node123" return daemon @@ -532,29 +492,3 @@ async def test_a_burst_of_changes_produces_one_broadcast(tmp_path, shared_dir, g await asyncio.sleep(0.05) session._send.assert_called_once() - - -@pytest.mark.asyncio -async def test_swarm_registration_only_sends_new_hashes_after_the_first( - tmp_path, shared_dir, gek): - daemon = _new_daemon_for_group(tmp_path, shared_dir, gek, visibility="public") - indexer = DirectoryIndexer( - roots=one_root(shared_dir), group_id="a" * 32, - sk_node=Ed25519PrivateKey.generate(), gek=gek) - await indexer.initial_scan() - total_files = indexer.index.count - - await daemon._on_index_change(indexer) - await asyncio.sleep(0.05) - assert len(daemon._hub.register_swarm.call_args_list[0].args[0]) == total_files - - from meshbay_common.protocol import IndexEntry - indexer.index.add_entry(IndexEntry(id="new-file-id", name="new.mp4", - path="shared", size=10, type="video", - added_at=0)) - await daemon._on_index_change(indexer) - await asyncio.sleep(0.05) - - assert daemon._hub.register_swarm.call_count == 2 - assert daemon._hub.register_swarm.call_args_list[1].args[0] == ["new-file-id"], \ - "only the newly added hash must be (re-)registered, not the whole library" diff --git a/packages/meshbay-node/tests/test_indexer.py b/packages/meshbay-node/tests/test_indexer.py index e97e2ef..c60600f 100644 --- a/packages/meshbay-node/tests/test_indexer.py +++ b/packages/meshbay-node/tests/test_indexer.py @@ -42,58 +42,6 @@ def shared_dir(tmp_path): # ── GroupIndex tests ────────────────────────────────────────────────────────── -def test_group_index_serialize_deserialize_private(sk_node, gek, shared_dir): - idx = GroupIndex(group_id="grp-001", sk_node=sk_node, gek=gek) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry( - id="abc123", name="video.mkv", path="", size=1024, - type="video", added_at=int(time.time()), duration=120)) - - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - - assert recovered.group_id == "grp-001" - assert recovered.count == 1 - assert recovered.entries[0].name == "video.mkv" - assert recovered.entries[0].type == "video" - - -def test_group_index_serialize_deserialize_public(sk_node): - idx = GroupIndex(group_id="pub-001", sk_node=sk_node, gek=None) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry( - id="xyz789", name="readme.txt", path="", size=42, - type="document", added_at=int(time.time()))) - - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=None) - assert recovered.count == 1 - assert recovered.entries[0].id == "xyz789" - - -def test_group_index_wrong_gek_rejected(sk_node, gek): - idx = GroupIndex(group_id="grp-002", sk_node=sk_node, gek=gek) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry(id="a", name="f.mp3", path="", size=1, - type="audio", added_at=0)) - wire = idx.serialize() - - wrong_gek = generate_gek() - with pytest.raises(Exception): # InvalidTag from AEAD - GroupIndex.deserialize(wire, sk_node=sk_node, gek=wrong_gek) - - -def test_group_index_tampered_rejected(sk_node, gek): - idx = GroupIndex(group_id="grp-003", sk_node=sk_node, gek=gek) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry(id="b", name="f.mp4", path="", size=1, - type="video", added_at=0)) - wire = bytearray(idx.serialize()) - wire[-5] ^= 0xFF # flip bytes at the end - with pytest.raises(Exception): - GroupIndex.deserialize(bytes(wire), sk_node=sk_node, gek=gek) - - def test_group_index_diff(sk_node, gek): from meshbay_common.protocol import IndexEntry v1 = GroupIndex(group_id="g", sk_node=sk_node, gek=gek, version=1) @@ -229,16 +177,6 @@ async def test_on_change_callback(shared_dir, sk_node, gek): assert len(changes) >= 1, "on_change should have been called" -@pytest.mark.asyncio -async def test_index_roundtrip_after_scan(shared_dir, sk_node, gek): - indexer = DirectoryIndexer(roots=one_root(shared_dir), group_id="g", sk_node=sk_node, gek=gek) - await indexer.initial_scan() - - wire = indexer.index.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - assert recovered.count == indexer.index.count - - # ── Cache-aware scanning ─────────────────────────────────────────────────────── @pytest.fixture @@ -796,35 +734,13 @@ def test_partial_hash_is_deterministic(tmp_path): assert e1.id == e2.id -def test_group_index_roundtrip_preserves_hash_version(sk_node, gek): - from meshbay_common.protocol import IndexEntry - idx = GroupIndex(group_id="hv-test", sk_node=sk_node, gek=gek) - idx.add_entry(IndexEntry( - id="aaa", name="small.mp4", path="root", size=1024, - type="video", added_at=100, hash_version=1)) - idx.add_entry(IndexEntry( - id="bbb", name="big.mkv", path="root", size=50_000_000, - type="video", added_at=200, hash_version=2)) - - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - - by_id = {e.id: e for e in recovered.entries} - assert by_id["aaa"].hash_version == 1 - assert by_id["bbb"].hash_version == 2 - - -def test_deserialize_without_hash_version_defaults_to_1(sk_node, gek): - """Entries serialized by old code (no hash_version field) must deserialize - as hash_version=1.""" +def test_an_entry_without_hash_version_defaults_to_1(): + """An entry built without the field (written before it existed) reads as a + full-read hash.""" from meshbay_common.protocol import IndexEntry - idx = GroupIndex(group_id="compat", sk_node=sk_node, gek=gek) - idx.add_entry(IndexEntry( - id="old", name="f.mp4", path="root", size=1024, - type="video", added_at=100)) - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - assert recovered.entries[0].hash_version == 1 + e = IndexEntry(id="old", name="f.mp4", path="root", size=1024, + type="video", added_at=100) + assert e.hash_version == 1 def test_index_entry_wire_includes_hash_version(): diff --git a/packages/meshbay-node/tests/test_security_regressions.py b/packages/meshbay-node/tests/test_security_regressions.py index c0c2c2c..43e88c2 100644 --- a/packages/meshbay-node/tests/test_security_regressions.py +++ b/packages/meshbay-node/tests/test_security_regressions.py @@ -602,19 +602,18 @@ def test_denylist_persists_and_honours_groups(tmp_path): assert not reloaded.is_denied("someone", "other", "g-allowed") -def test_swarm_registration_skips_private_groups(): +def test_the_node_registers_no_content_hash_with_the_hub(): """ - H7: the daemon registered content hashes for every group, private included, - handing the hub a fingerprint of every private file. The bug was masked by a - mis-mounted route, so fixing the route without this filter would have turned a - dormant leak into a live one. + H7: the daemon once registered content hashes for every group, private + included, handing the hub a fingerprint of every private file. The swarm that + received them is gone, so no group's hashes are sent to the hub at all. """ - source = daemon_source() - assert 'visibility' in source and '_register_swarm' in source - # Both registration sites must gate on public visibility. - for marker in ['gctx.get("visibility") == "public"', - 'group_cfg.visibility == "public"']: - assert marker in source, f"swarm registration not gated: {marker}" + from pathlib import Path + + import meshbay_node.hub_client as hub_client + for source in (daemon_source(), Path(hub_client.__file__).read_text(encoding="utf-8")): + assert "/v1/swarm" not in source + assert "register_swarm" not in source def test_keystore_argon2_is_production_strength(): diff --git a/packages/meshbay-node/tests/test_tmdb_config_policy.py b/packages/meshbay-node/tests/test_tmdb_config_policy.py index 6a51eb0..c671ac8 100644 --- a/packages/meshbay-node/tests/test_tmdb_config_policy.py +++ b/packages/meshbay-node/tests/test_tmdb_config_policy.py @@ -22,7 +22,7 @@ from pathlib import Path import pytest from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey -from meshbay_common.adminop import OP_TMDB_CONFIG +from meshbay_common.adminop import OP_TMDB_CONFIG, tmdb_config_subject from meshbay_node.indexer.group_index import GroupIndex from meshbay_node.roster import Roster from meshbay_node.transport.webrtc_server import WebRTCPeerSession @@ -134,7 +134,9 @@ async def test_subject_reflects_whether_a_token_was_supplied(tmp_path): session._do_tmdb_config({"token": "x"}) _, subject, _, _ = issued[0] - assert subject == "custom_token=yes,language=default" + # The token is bound by its digest and never written into the subject. + assert subject == tmdb_config_subject("x", None) + assert "sha256:" in subject and tmdb_config_subject("y", None) != subject async def test_subject_says_no_custom_token_when_none_given(tmp_path): @@ -146,7 +148,10 @@ async def test_subject_says_no_custom_token_when_none_given(tmp_path): session._do_tmdb_config({}) _, subject, _, _ = issued[0] - assert subject == "custom_token=no,language=default" + # Nothing given means both unchanged — distinct from clearing either. + assert subject == tmdb_config_subject(None, None) + assert subject != tmdb_config_subject("", None) + assert subject != tmdb_config_subject(None, "") async def test_subject_reflects_a_configured_language(tmp_path): @@ -158,7 +163,7 @@ async def test_subject_reflects_a_configured_language(tmp_path): session._do_tmdb_config({"language": "fr-FR"}) _, subject, payload, _ = issued[0] - assert subject == "custom_token=no,language=fr-FR" + assert subject == tmdb_config_subject(None, "fr-FR") assert payload["language"] == "fr-FR" diff --git a/packages/meshbay-node/tests/test_webrtc_transport.py b/packages/meshbay-node/tests/test_webrtc_transport.py index 990b1da..7a6f517 100644 --- a/packages/meshbay-node/tests/test_webrtc_transport.py +++ b/packages/meshbay-node/tests/test_webrtc_transport.py @@ -34,13 +34,12 @@ from meshbay_common.adminop import ( OP_INVITE_CREATE, OP_INVITE_LINK_CREATE, admin_transcript, + invite_create_subject, ) from meshbay_common.crypto import ( generate_gek, pk_to_b64, - unwrap_gek, unwrap_gek_aes, - wrap_gek, wrap_gek_aes, ) from meshbay_common.groupbox import PURPOSE_ACK, PURPOSE_INDEX, unseal @@ -1337,7 +1336,8 @@ async def test_invite_then_join_delivers_the_gek(sk_node, sk_hub, gek, shared_di challenge_msg = await asyncio.wait_for(q_admin.get(), timeout=5.0) assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE assert challenge_msg["op"] == OP_INVITE_CREATE - assert challenge_msg["subject"] == "user-002" + # The name the invitation records is signed with the account it is for. + assert challenge_msg["subject"] == invite_create_subject("user-002", "bob") ch_admin.send(_pack({ "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, @@ -1570,7 +1570,7 @@ async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_di await bundle_store.open() # Pre-populate a bundle for user-001 in group "g" - bundle = wrap_gek(gek, pk_x_raw) + bundle = wrap_gek_aes(gek, pk_x_raw) await bundle_store.store("g", "user-001", bundle["pk_eph_b64"], bundle["nonce_b64"], bundle["wrapped_b64"]) @@ -1631,7 +1631,7 @@ async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_di assert bundle_resp["found"] is True # Step 3: Unwrap GEK and compute HMAC proof - recovered_gek = unwrap_gek(bundle_resp, sk_x_raw, pk_x_raw) + recovered_gek = unwrap_gek_aes(bundle_resp, sk_x_raw, pk_x_raw) assert recovered_gek == gek nonce_s = base64.b64decode(msg["nonce"]) diff --git a/packages/meshbay-node/tests/transfer_probe.py b/packages/meshbay-node/tests/transfer_probe.py index 7aa4902..6666960 100755 --- a/packages/meshbay-node/tests/transfer_probe.py +++ b/packages/meshbay-node/tests/transfer_probe.py @@ -274,7 +274,14 @@ async def operator_checks(client, ack, group, node_id) -> int: import re as _re failures = 0 - member_cap = int((ack.get("transfer_limits") or {}).get("download") or 0) + # This member's own cap, as the node states it on every `transfer_state`. + probe_tr, probe_state = await _open_transfer(client) + member_cap = int(probe_state.get("cap") or 0) + client.send({"type": "transfer_close", "v": "0.1", "tr": probe_tr, "reason": "done"}) + while True: # closed before anything below counts what is in use + reply = await client.recv_type("transfer_state", timeout=15) + if reply.get("tr") == probe_tr and reply.get("state") == "closed": + break # The queue has to be held by the *node* cap, not by this member's own. # @@ -446,16 +453,6 @@ async def probe(args) -> int: ack = await alice.connect(http, group["id"], node_id) - limits = ack.get("transfer_limits") - if limits is None: - print("This node does not hand out transfer slots — it predates " - "them, or the handshake ack lost the field. Nothing below " - "can be measured.") - await alice.close() - return 1 - cap = int(limits.get("download") or 0) - print(f"node reports this member may run {cap} download(s) at once\n") - if args.operator: return await operator_checks(alice, ack, group, node_id) @@ -472,9 +469,17 @@ async def probe(args) -> int: # held. So the first reply is read for what the node says is already # in use, and the run stops rather than measuring against a moving # floor. - want = args.want or (cap + 2) opened = [await _open_transfer(alice)] first = opened[0][1] + # This member's cap, as the node states it on every `transfer_state`. + cap = int(first.get("cap") or 0) + if not cap: + print("This node does not state a transfer cap — nothing below can " + "be measured.") + await alice.close() + return 1 + print(f"node reports this member may run {cap} download(s) at once\n") + want = args.want or (cap + 2) if first.get("state") != "granted" or first.get("used", 1) != 1: print(f"this node is not idle: it reports {first.get('used')} of " f"{first.get('cap')} slots already used by this member, and " |