diff options
Diffstat (limited to 'packages/meshbay-node/tests')
| -rw-r--r-- | packages/meshbay-node/tests/golden/cli.json | 9 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_attach_from_the_hub.py | 99 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_chat_is_bounded.py | 35 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_ffprobe_is_bounded.py | 61 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_linkpreview.py | 94 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_member_capacity.py | 151 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_member_errors_are_plain.py | 22 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_ops.py | 2 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_partial_uploads.py | 118 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_revocation_closes_sessions.py | 53 |
10 files changed, 621 insertions, 23 deletions
diff --git a/packages/meshbay-node/tests/golden/cli.json b/packages/meshbay-node/tests/golden/cli.json index b397eeb..48d2f74 100644 --- a/packages/meshbay-node/tests/golden/cli.json +++ b/packages/meshbay-node/tests/golden/cli.json @@ -4,7 +4,7 @@ "asked": [], "exit": 0, "stderr": "", - "stdout": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\n\nMeshBay Node daemon\n\npositional arguments:\n {init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}\n init: provision config + keystore | reset: erase all node state | status:\n node state and keys | operator pair: pair a browser with this node |\n member list|invite|cancel|revoke|unpin | group list|add|remove | root\n list|add|remove|set|eject|plug | gek init|rotate | file list|rm | video\n rematch: re-resolve TMDB matches for a group's videos | chat\n status|rotate|encrypt-history|prune | denylist show|clear | stun\n list|add|remove|reset | transfers show|set|max-size|per-member: live\n transfer slots, the node-wide caps, the largest single upload, and how\n many one member may run at once in a group | reload: re-read node.toml\n (hot; systemd or the loopback API) | restart-daemon: restart the node\n (systemd unit, the Windows autostart launcher, or the service task,\n whichever applies) | autostart install|remove|start|stop|status (Windows:\n run meshbay-node at each sign-in, no admin) | service\n install|remove|start|stop|status (Windows: run at boot, before sign-in,\n needs admin once to install) | calibrate-argon2: benchmark\n subcommand 'pair' for operator; list|invite|revoke|unpin for member; list|add|remove\n for group; list|add|remove|set|eject|plug for root; init|rotate for gek;\n list|rm for file; rematch for video; show|clear for denylist;\n list|add|remove|reset for stun; show|set|max-size|per-member for\n transfers; install|remove|start|stop|status for autostart and for service\n target username for member invite|revoke|unpin (an optional e-mail label with\n --link, a link id for member cancel); group name for group add; file id\n for file rm; identifier for denylist clear; download cap for transfers\n set; size in GB for transfers max-size\n value the second value where a verb takes two: the upload cap for transfers set\n\noptions:\n -h, --help show this help message and exit\n --hub-url HUB_URL hub URL, for init (e.g. https://meshbay.org)\n --username USERNAME hub username, for init\n --dir DIR shared directory, for group add\n --yes skip the confirmation for destructive commands\n --config CONFIG Config file path\n --group GROUP group id (optional if only one is configured)\n --link member invite: an invitation link, for someone who may have no account yet\n (valid 7 days, single use)\n --writable root accepts member uploads (root add/set)\n --no-writable root is read-only (root add/set, group add)\n --removable mark root as removable (root set/add)\n --no-removable mark root as not removable (root set)\n --name NAME root name (root add; defaults to directory basename)\n --log-level {DEBUG,INFO,WARNING,ERROR}\n", + "stdout": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--open] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\n\nMeshBay Node daemon\n\npositional arguments:\n {init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}\n init: provision config + keystore | reset: erase all node state | status:\n node state and keys | operator pair: pair a browser with this node |\n member list|invite|cancel|revoke|unpin | group list|add|remove | root\n list|add|remove|set|eject|plug | gek init|rotate | file list|rm | video\n rematch: re-resolve TMDB matches for a group's videos | chat\n status|rotate|encrypt-history|prune | denylist show|clear | stun\n list|add|remove|reset | transfers show|set|max-size|per-member: live\n transfer slots, the node-wide caps, the largest single upload, and how\n many one member may run at once in a group | reload: re-read node.toml\n (hot; systemd or the loopback API) | restart-daemon: restart the node\n (systemd unit, the Windows autostart launcher, or the service task,\n whichever applies) | autostart install|remove|start|stop|status (Windows:\n run meshbay-node at each sign-in, no admin) | service\n install|remove|start|stop|status (Windows: run at boot, before sign-in,\n needs admin once to install) | calibrate-argon2: benchmark\n subcommand 'pair' for operator; list|invite|revoke|unpin for member; list|add|remove\n for group; list|add|remove|set|eject|plug for root; init|rotate for gek;\n list|rm for file; rematch for video; show|clear for denylist;\n list|add|remove|reset for stun; show|set|max-size|per-member for\n transfers; install|remove|start|stop|status for autostart and for service\n target username for member invite|revoke|unpin (an optional e-mail label with\n --link, a link id for member cancel); group name for group add; file id\n for file rm; identifier for denylist clear; download cap for transfers\n set; size in GB for transfers max-size\n value the second value where a verb takes two: the upload cap for transfers set\n\noptions:\n -h, --help show this help message and exit\n --hub-url HUB_URL hub URL, for init (e.g. https://meshbay.org)\n --username USERNAME hub username, for init\n --dir DIR shared directory, for group add\n --yes skip the confirmation for destructive commands\n --config CONFIG Config file path\n --group GROUP group id (optional if only one is configured)\n --link member invite: an invitation link, for someone who may have no account yet\n (valid 7 days, single use)\n --writable root accepts member uploads (root add/set)\n --no-writable root is read-only (root add/set, group add)\n --removable mark root as removable (root set/add)\n --no-removable mark root as not removable (root set)\n --open group add: anyone the hub lists the group to may join (default: by\n invitation only)\n --name NAME root name (root add; defaults to directory basename)\n --log-level {DEBUG,INFO,WARNING,ERROR}\n", "systemctl": [] }, "autostart no-such-sub": { @@ -266,7 +266,7 @@ "asked": [], "exit": 1, "stderr": "", - "stdout": "usage: meshbay-node group add <name> --dir <path> [--no-writable]\n\nThe group must already exist on the hub and be yours. This\nonly tells the node to host it, and picks its first\ndirectory, which accepts uploads unless --no-writable.\nAdd more with: meshbay-node root add <path> [--writable]\n", + "stdout": "usage: meshbay-node group add <name> --dir <path> [--no-writable] [--open]\n\nThe group must already exist on the hub and be yours. This\nonly tells the node to host it, and picks its first\ndirectory, which accepts uploads unless --no-writable.\nAdd more with: meshbay-node root add <path> [--writable]\n", "systemctl": [] }, "group add g --dir /tmp/media --no-writable": { @@ -275,6 +275,7 @@ "POST", "/api/groups/attach", { + "join_policy": "invite", "name": "g", "shared_dir": "/tmp/media", "writable": false @@ -284,7 +285,7 @@ "asked": [], "exit": 0, "stderr": "", - "stdout": "g (g) added to <tmp>/node.toml\n shared_dir <tmp> (read-only)\n\nTell the daemon to re-read its config, then give the group a key:\n meshbay-node reload\n meshbay-node gek init --group g\n\nThe key is this group's own — members of your other groups cannot\nread it, and joining one says nothing about the other.\n", + "stdout": "g (g) added to <tmp>/node.toml\n shared_dir <tmp> (read-only)\n join_policy invite\n\nTell the daemon to re-read its config, then give the group a key:\n meshbay-node reload\n meshbay-node gek init --group g\n\nThe key is this group's own — members of your other groups cannot\nread it, and joining one says nothing about the other.\n", "systemctl": [] }, "group list": { @@ -423,7 +424,7 @@ "api": [], "asked": [], "exit": 2, - "stderr": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\nmeshbay-node: error: argument command: invalid choice: 'no-such-verb' (choose from init, reset, status, gek-init, gek, operator, member, group, root, file, video, chat, denylist, stun, transfers, reload, restart-daemon, autostart, service, calibrate-argon2)\n", + "stderr": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--open] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\nmeshbay-node: error: argument command: invalid choice: 'no-such-verb' (choose from init, reset, status, gek-init, gek, operator, member, group, root, file, video, chat, denylist, stun, transfers, reload, restart-daemon, autostart, service, calibrate-argon2)\n", "stdout": "", "systemctl": [] }, diff --git a/packages/meshbay-node/tests/test_attach_from_the_hub.py b/packages/meshbay-node/tests/test_attach_from_the_hub.py new file mode 100644 index 0000000..f4aaa25 --- /dev/null +++ b/packages/meshbay-node/tests/test_attach_from_the_hub.py @@ -0,0 +1,99 @@ +""" +Hosting a group writes what the operator said, never what the hub says. + +`attach_group` looks the group up on the hub, because the node is the process +signed in there. What it must not take from that answer is how people join: a +hub able to declare a group open would be handed its key by anyone it sent +(admission in `transport/webrtc/admission.py` admits a stranger to an open +group). And every string it writes into node.toml is someone else's text — a +group name chosen on the hub, a folder name — so none of it may end a TOML +string and write lines of its own. +""" + +import tomllib +from pathlib import Path + +import pytest +from meshbay_node import ops +from meshbay_node.config import load_config +from meshbay_node.ops.node_toml import toml_string + +GID = "0f8fad5b-d9cb-469f-a165-70867728950e" +HOSTILE = 'Films"\n[node]\nui_port = 1\n# ' + + +class _Hub: + _session = object() # signed in + + def __init__(self, group): + self._group = group + + async def list_my_groups(self): + return [self._group] + + +def _state(tmp_path: Path, group: dict) -> dict: + conf = tmp_path / "node.toml" + conf.write_text('[hub]\nurl = "https://hub.invalid"\nusername = "op"\n\n' + '[node]\nui_port = 18000\n', encoding="utf-8") + return {"config": load_config(conf), "config_path": str(conf), "hub": _Hub(group)} + + +def _hosted(tmp_path: Path) -> dict: + parsed = tomllib.loads((tmp_path / "node.toml").read_text(encoding="utf-8")) + return parsed["groups"][0] | {"node": parsed["node"]} + + +async def test_a_group_the_hub_calls_open_is_hosted_by_invitation(tmp_path): + state = _state(tmp_path, {"id": GID, "name": "Films", "visibility": "public", + "join_policy": "open"}) + out = await ops.attach_group(state, "Films", str(tmp_path / "share")) + + hosted = _hosted(tmp_path) + assert hosted["join_policy"] == "invite" + assert hosted["visibility"] == "private" + assert out["hub_join_policy"] == "open", "the caller is not told the two differ" + + +async def test_the_operator_opens_it(tmp_path): + state = _state(tmp_path, {"id": GID, "name": "Films", "join_policy": "invite"}) + await ops.attach_group(state, "Films", str(tmp_path / "share"), join_policy="open") + + hosted = _hosted(tmp_path) + assert (hosted["join_policy"], hosted["visibility"]) == ("open", "public") + + +async def test_an_unknown_policy_is_refused(tmp_path): + state = _state(tmp_path, {"id": GID, "name": "Films"}) + with pytest.raises(ops.OpError): + await ops.attach_group(state, "Films", str(tmp_path / "share"), + join_policy="anyone") + + +async def test_a_group_name_cannot_write_lines_into_node_toml(tmp_path): + state = _state(tmp_path, {"id": GID, "name": HOSTILE}) + await ops.attach_group(state, HOSTILE, str(tmp_path / "share")) + + hosted = _hosted(tmp_path) + assert hosted["name"] == HOSTILE + assert hosted["node"]["ui_port"] == 18000 + + +async def test_a_folder_name_cannot_either(tmp_path): + state = _state(tmp_path, {"id": GID, "name": "Films"}) + await ops.attach_group(state, "Films", str(tmp_path / "share")) + state["config"] = load_config(tmp_path / "node.toml") + # Root names are refused with such characters already (roots.py); the + # path is not, and is the operator's own folder or a member's request. + weird = tmp_path / 'a "quoted"\\ folder\n[node]' + await ops.add_root(state, GID, str(weird), name="extra") + + hosted = _hosted(tmp_path) + assert Path(hosted["roots"][-1]["path"]) == weird + assert hosted["node"]["ui_port"] == 18000 + + +@pytest.mark.parametrize("value", ["plain", 'q"uote', "back\\slash", "line\nbreak", + "tab\there", "del\x7f", "café 日本"]) +def test_every_string_reads_back_as_written(value): + assert tomllib.loads(f"v = {toml_string(value)}")["v"] == value diff --git a/packages/meshbay-node/tests/test_chat_is_bounded.py b/packages/meshbay-node/tests/test_chat_is_bounded.py index 3332601..c48f26d 100644 --- a/packages/meshbay-node/tests/test_chat_is_bounded.py +++ b/packages/meshbay-node/tests/test_chat_is_bounded.py @@ -204,3 +204,38 @@ async def test_one_member_at_their_limit_has_not_spent_anyone_elses(ctx, store): await _flood(bob, ctx, 1) assert not _errors(bob), "one member's flood silenced another" assert _acks(bob) + + +# ── the fields beside the ciphertext ───────────────────────────────────────── + +async def test_the_clear_fields_are_what_they_claim_and_no_larger(ctx, store): + """ + `sender_name` and `thread_id` travel in clear beside the sealed envelope, + which carries its own. They are stored on the operator's disk and relayed + to every member, so a megabyte of name or a list for a thread id is dropped, + not kept. + """ + alice = _session(ctx, store, user="alice", conn="c1") + bob = _session(ctx, store, user="bob", conn="c2") + msg = _message(alice, size=64) + msg["sender_name"] = "x" * (1024 * 1024) + msg["thread_id"] = list(range(10_000)) + alice._do_chat_message(msg) + await _drain(ctx) + + stored = (await store.get_recent(limit=1))[0] + assert stored.sender_name in ("", None) + assert stored.thread_id is None + relayed = [m for m in bob.sent if m.get("type") == "chat_msg"] + assert relayed and relayed[-1]["sender_name"] == "" and relayed[-1]["thread_id"] is None + + +async def test_ordinary_clear_fields_pass_unchanged(ctx, store): + alice = _session(ctx, store, user="alice", conn="c1") + msg = _message(alice, size=64) + msg["sender_name"] = "Alice" + msg["thread_id"] = "42" + alice._do_chat_message(msg) + await _drain(ctx) + stored = (await store.get_recent(limit=1))[0] + assert (stored.sender_name, stored.thread_id) == ("Alice", "42") diff --git a/packages/meshbay-node/tests/test_ffprobe_is_bounded.py b/packages/meshbay-node/tests/test_ffprobe_is_bounded.py new file mode 100644 index 0000000..c96a19a --- /dev/null +++ b/packages/meshbay-node/tests/test_ffprobe_is_bounded.py @@ -0,0 +1,61 @@ +""" +ffprobe over a member's file is bounded, and a bounded wait stops the process. + +`probe_video` runs before every stream and subtitle request and in the +enrichment pool. A file that keeps ffprobe busy must not hold any of them, and +giving up on the wait is not enough: an asyncio subprocess whose wait was +cancelled keeps running. These run a stand-in ffprobe that never answers and +check both — the call returns, and the process is gone. +""" + +import asyncio +import os +import sys +import time + +import pytest +from meshbay_node import media_probe, platform + +pytestmark = pytest.mark.skipif(sys.platform == "win32", reason="a POSIX shell stand-in") + + +@pytest.fixture +def hanging_ffprobe(tmp_path, monkeypatch): + pid_file = tmp_path / "pid" + tool = tmp_path / "ffprobe" + tool.write_text(f"#!/bin/sh\necho $$ > {pid_file}\nexec sleep 600\n") + tool.chmod(0o755) + monkeypatch.setattr(platform, "_ffprobe_path", str(tool)) + return pid_file + + +def _gone(pid: int) -> bool: + try: + os.kill(pid, 0) + except ProcessLookupError: + return True + # A zombie still answers kill(0); its state says it has exited. + try: + with open(f"/proc/{pid}/stat") as f: + return f.read().split()[2] == "Z" + except OSError: + return True + + +async def test_a_probe_that_never_answers_times_out_and_is_stopped(hanging_ffprobe, monkeypatch): + monkeypatch.setattr(media_probe, "FFPROBE_TIMEOUT_SECS", 0.5) + started = time.monotonic() + with pytest.raises(RuntimeError, match="timed out"): + await media_probe.probe_video("/nonexistent/file.mkv") + assert time.monotonic() - started < 5 + assert _gone(int(hanging_ffprobe.read_text())) + + +async def test_a_caller_giving_up_stops_it_too(hanging_ffprobe): + with pytest.raises(TimeoutError): + await asyncio.wait_for(media_probe.probe_video("/nonexistent/file.mkv"), 0.5) + for _ in range(50): + if hanging_ffprobe.exists(): + break + await asyncio.sleep(0.05) + assert _gone(int(hanging_ffprobe.read_text())) diff --git a/packages/meshbay-node/tests/test_linkpreview.py b/packages/meshbay-node/tests/test_linkpreview.py index 9fca186..3e6eaf7 100644 --- a/packages/meshbay-node/tests/test_linkpreview.py +++ b/packages/meshbay-node/tests/test_linkpreview.py @@ -6,12 +6,13 @@ decides an outbound request from the operator's machine. Anything that is not a public http(s) address must be refused before a socket opens. """ +import asyncio import socket import httpx import pytest from meshbay_node import linkpreview -from meshbay_node.linkpreview import UnsafeURL, safe_url +from meshbay_node.linkpreview import UnsafeURL, check_url, safe_url PUBLIC_IP = "93.184.216.34" # example.com, historically @@ -43,9 +44,9 @@ def resolves_public(monkeypatch): "javascript:alert(1)", "not a url", ]) -def test_safe_url_refuses(url): +async def test_check_url_refuses(url): with pytest.raises(UnsafeURL): - safe_url(url) + await check_url(url) @pytest.mark.parametrize("url", [ @@ -72,11 +73,11 @@ def test_safe_url_allows_the_web_ports(url, resolves_public): assert safe_url(url) == url -def test_safe_url_accepts_a_public_host(resolves_public): - assert safe_url("https://example.com/some/page") == "https://example.com/some/page" +async def test_check_url_accepts_a_public_host(resolves_public): + assert await check_url("https://example.com/some/page") == "https://example.com/some/page" -def test_safe_url_refuses_a_host_with_any_private_record(monkeypatch): +async def test_check_url_refuses_a_host_with_any_private_record(monkeypatch): def mixed(host, port, *a, **k): return [ (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, "", (PUBLIC_IP, port)), @@ -84,7 +85,7 @@ def test_safe_url_refuses_a_host_with_any_private_record(monkeypatch): ] monkeypatch.setattr(linkpreview.socket, "getaddrinfo", mixed) with pytest.raises(UnsafeURL): - safe_url("https://sneaky.example/x") + await check_url("https://sneaky.example/x") # ── fetch_preview ────────────────────────────────────────────────────────── @@ -273,3 +274,82 @@ async def test_a_declared_oversized_image_is_not_read(resolves_public): async with _client(handler) as c: assert await linkpreview.fetch_image("https://example.com/x.png", client=c) is None assert counter["sent"] <= _CHUNK + + +# ── The connection goes where the check said ─────────────────────────────── + +class _Recorder: + """Stands in for the real socket layer under the pinned backend.""" + def __init__(self): + self.hosts = [] + + async def connect_tcp(self, host, port, **kw): + self.hosts.append(host) + raise httpx.ConnectError("recorded, not connected") + + +def _pinned_with(recorder): + backend = linkpreview._PinnedBackend() + backend._inner = recorder + return backend + + +async def test_the_socket_is_opened_to_the_checked_address(resolves_public): + rec = _Recorder() + with pytest.raises(httpx.ConnectError): + await _pinned_with(rec).connect_tcp("example.com", 443) + assert rec.hosts == [PUBLIC_IP], "the name, not the checked address, was dialled" + + +async def test_a_name_that_rebinds_never_reaches_the_lan(monkeypatch): + """ + Answers clean when checked, then with a LAN address. Checked once and + dialled by name, the request would go to the LAN before anything looked; + resolved and checked by the backend that dials, it goes nowhere. + """ + answers = iter([PUBLIC_IP, "192.168.1.1", "192.168.1.1"]) + + def rebinding(host, port, *a, **k): + return [(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, "", + (next(answers), port))] + monkeypatch.setattr(linkpreview.socket, "getaddrinfo", rebinding) + + await check_url("http://rebind.example/x") # the clean answer + rec = _Recorder() + with pytest.raises(UnsafeURL): + await _pinned_with(rec).connect_tcp("rebind.example", 80) + assert rec.hosts == [] + + +async def test_the_real_client_is_pinned(): + """What `fetch_preview` uses when the caller gives no client.""" + client = linkpreview._new_client() + try: + pool = client._transport._pool + assert isinstance(pool._network_backend, linkpreview._PinnedBackend) + assert client._trust_env is False, "a proxy from the environment would unpin it" + finally: + await client.aclose() + + +async def test_resolving_does_not_hold_the_event_loop(monkeypatch): + import time as _time + + def slow(host, port, *a, **k): + _time.sleep(0.4) + return [(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, "", (PUBLIC_IP, port))] + monkeypatch.setattr(linkpreview.socket, "getaddrinfo", slow) + + ticks = 0 + + async def ticker(): + nonlocal ticks + while True: + await asyncio.sleep(0.02) + ticks += 1 + t = asyncio.create_task(ticker()) + try: + await check_url("https://slow.example/x") + finally: + t.cancel() + assert ticks >= 10, "the event loop stood still while a name resolved" diff --git a/packages/meshbay-node/tests/test_member_capacity.py b/packages/meshbay-node/tests/test_member_capacity.py new file mode 100644 index 0000000..16d112d --- /dev/null +++ b/packages/meshbay-node/tests/test_member_capacity.py @@ -0,0 +1,151 @@ +""" +What one member may hold of a node: a share, sized so that real use never +meets it. + +The heaviest real member — twenty groups on this node, three devices and a +spare tab — holds up to 52 peer sessions (twelve per device for search and +music, plus the open group page). The node holds 128, and one account at most +half. One account may play half the node's video slots and run two subtitle +extractions. The node's own account is the operator's machine and is not +counted. And a message is at most 8 MiB, decoded with a bound on every +container, because one message of tiny elements decodes to many times its size. +""" + +import struct +from unittest.mock import MagicMock + +import msgpack +import pytest +from meshbay_node.transport.webrtc.channel import _DataChannelBuffer +from meshbay_node.transport.webrtc.limits import MAX_MSG +from meshbay_node.transport.webrtc_server import ( + MAX_PEER_SESSIONS, + MAX_PEER_SESSIONS_PER_ACCOUNT, + WebRTCPeerSession, + WebRTCTransport, +) + + +class _Held: + def __init__(self, user): + self._offer_user = user + self._user_id = user + + async def close(self): + pass + + +def _transport(**ctx) -> WebRTCTransport: + tp = WebRTCTransport(sk_node=MagicMock(), hub_pk_pem=b"", gek=None, + roots=None, index=None) + tp._ctx.update(ctx) + return tp + + +def test_the_shares_are_what_was_agreed(): + assert MAX_PEER_SESSIONS == 128 + assert MAX_PEER_SESSIONS_PER_ACCOUNT == 64 + assert MAX_MSG == 8 * 1024 * 1024 + # The heaviest real member fits with room to spare. + assert 4 * (12 + 1) < MAX_PEER_SESSIONS_PER_ACCOUNT + + +@pytest.mark.asyncio +async def test_one_account_cannot_hold_more_than_its_share(): + tp = _transport() + for i in range(MAX_PEER_SESSIONS_PER_ACCOUNT): + tp._sessions[f"a-{i}"] = _Held("alice") + + with pytest.raises(RuntimeError, match="share"): + await tp.handle_offer("v=0", "alice-one-more", "alice") + assert "alice-one-more" not in tp._sessions + + # Somebody else still gets in: what stops this offer is not the share. + try: + await tp.handle_offer("v=0", "bob-first", "bob") + except RuntimeError as e: + assert "share" not in str(e) and "limit" not in str(e) + except Exception: + pass + + +@pytest.mark.asyncio +async def test_the_operators_own_account_is_not_counted(): + tp = _transport(node_user_id="operator") + for i in range(MAX_PEER_SESSIONS_PER_ACCOUNT): + tp._sessions[f"o-{i}"] = _Held("operator") + try: + await tp.handle_offer("v=0", "operator-more", "operator") + except RuntimeError as e: + assert "share" not in str(e) + except Exception: + pass + + +def _session(ctx, user): + s = WebRTCPeerSession.__new__(WebRTCPeerSession) + s._ctx = ctx + s._user_id = user + s.sent = [] + s._send = s.sent.append + return s + + +def test_an_accounts_devices_share_one_allowance(): + ctx = {} + phone, laptop = _session(ctx, "alice"), _session(ctx, "alice") + with phone._account_share("streams", 2) as a, laptop._account_share("streams", 2) as b: + assert a and b + with phone._account_share("streams", 2) as c: + assert c is False + with _session(ctx, "bob")._account_share("streams", 2) as d: + assert d, "another member is not counted against alice" + with phone._account_share("streams", 2) as e: + assert e, "a place is given back when its stream ends" + + +def test_the_operator_is_not_counted_for_streams_either(): + ctx = {"node_user_id": "operator"} + s = _session(ctx, "operator") + with s._account_share("streams", 1) as a, s._account_share("streams", 1) as b: + assert a and b + + +@pytest.mark.asyncio +async def test_a_stream_past_the_accounts_share_is_refused_by_name(): + ctx = {"max_concurrent_streams": 8, "_streams_by_account": {"alice": 4}} + s = _session(ctx, "alice") + assert s._streams_per_account() == 4 + await s._stream_video({"file_id": "x"}) + assert s.sent[-1]["type"] == "error" + assert "Too many videos" in s.sent[-1]["detail"] + + +@pytest.mark.parametrize("cap,share", [(1, 1), (2, 1), (3, 2), (8, 4), (9, 5)]) +def test_the_stream_share_is_half_rounded_up(cap, share): + assert _session({"max_concurrent_streams": cap}, "a")._streams_per_account() == share + + +def _frame(obj) -> bytes: + body = msgpack.packb(obj, use_bin_type=True) + return struct.pack(">I", len(body)) + body + + +def test_the_largest_real_message_passes(): + buf = _DataChannelBuffer() + buf.feed(_frame({"type": "user_blob_store", "blob": b"x" * (1024 * 1024 + 64)})) + assert len(list(buf.messages())) == 1 + + +def test_a_message_of_tiny_elements_is_refused(): + buf = _DataChannelBuffer() + buf.feed(_frame({"type": "x", "items": [None] * 200_000})) + with pytest.raises(ValueError): + list(buf.messages()) + + +def test_a_message_over_the_limit_is_refused(): + buf = _DataChannelBuffer() + buf.feed(struct.pack(">I", MAX_MSG + 1) + b"x") + with pytest.raises(ValueError): + list(buf.messages()) diff --git a/packages/meshbay-node/tests/test_member_errors_are_plain.py b/packages/meshbay-node/tests/test_member_errors_are_plain.py new file mode 100644 index 0000000..e144e9b --- /dev/null +++ b/packages/meshbay-node/tests/test_member_errors_are_plain.py @@ -0,0 +1,22 @@ +""" +What a member is told when the media tools fail: that it failed. + +An exception's text from ffmpeg or ffprobe names the operator's paths, versions +and the libraries the build has; the member who asked needs none of it, and the +operator finds it in their log. Read from the source, because what matters is +that no reply in these handlers carries an exception's text at all. +""" + +import re +from pathlib import Path + +APPS = Path(__file__).resolve().parents[1] / "src" / "meshbay_node" / "transport" / "webrtc" + + +def test_no_reply_to_a_member_carries_an_exceptions_text(): + offenders = [] + for path in APPS.rglob("*.py"): + text = path.read_text(encoding="utf-8") + for m in re.finditer(r'"detail":\s*(f"[^"]*\{e\}[^"]*"|str\(e\))', text): + offenders.append(f"{path.name}: {m.group(0)}") + assert not offenders, offenders diff --git a/packages/meshbay-node/tests/test_ops.py b/packages/meshbay-node/tests/test_ops.py index 9eb8339..2c5a15a 100644 --- a/packages/meshbay-node/tests/test_ops.py +++ b/packages/meshbay-node/tests/test_ops.py @@ -405,7 +405,7 @@ def test_every_path_written_into_node_toml_goes_through_as_posix(): import re source = ops_source() # Every f-string interpolation that lands on the right of a TOML `path =`. - writes = re.findall(r'path\s*=\s*\\?"\{([^}]+)\}', source) + writes = re.findall(r'path\s*=\s*\\?"?\{(?:toml_string\()?([^}]+)\}', source) assert writes, "no TOML path writer found — did the config writer move?" for expr in writes: assert "as_posix()" in expr, ( diff --git a/packages/meshbay-node/tests/test_partial_uploads.py b/packages/meshbay-node/tests/test_partial_uploads.py index 80ab454..52aa306 100644 --- a/packages/meshbay-node/tests/test_partial_uploads.py +++ b/packages/meshbay-node/tests/test_partial_uploads.py @@ -17,10 +17,12 @@ may inherit — or overwrite the position of — the other's. """ import os +import re import time import types from pathlib import Path +import pytest from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from meshbay_common.crypto import generate_gek from meshbay_common.protocol import ( @@ -104,7 +106,7 @@ def _old(seconds: float) -> float: NOW = 1_000_000.0 -FILM = Path("/roots/media/film.mkv.part") +FILM = Path("/roots/media/film.mkv.0123abcd.part") def test_a_part_nobody_is_writing_and_nobody_has_touched_is_deleted(): @@ -138,7 +140,7 @@ def test_the_same_name_in_another_directory_does_not_protect_it(): path rebuilt from a root and a relative directory would be a second implementation that has to agree with the first for ever, and the state records the path it is writing instead.""" - other = Path("/roots/archive/film.mkv.part") + other = Path("/roots/archive/film.mkv.0123abcd.part") doomed = orphaned_parts([(other, _old(ORPHAN_AFTER_SECS + 1))], live={FILM}, now=NOW) assert doomed == [other] @@ -157,15 +159,15 @@ def test_a_finished_file_is_not_a_candidate(): def test_a_file_from_the_future_is_left_alone(): """A clock that went backwards is not evidence that a file is abandoned, and deleting is not reversible.""" - doomed = orphaned_parts([(Path("/roots/media/a.part"), NOW + 10_000)], + doomed = orphaned_parts([(Path("/roots/media/a.0123abcd.part"), NOW + 10_000)], live=set(), now=NOW) assert doomed == [] def test_the_boundary_is_the_age_itself(): - at = [(Path("/roots/media/a.part"), _old(ORPHAN_AFTER_SECS))] - just_under = [(Path("/roots/media/a.part"), _old(ORPHAN_AFTER_SECS - 1))] - assert orphaned_parts(at, set(), NOW) == [Path("/roots/media/a.part")] + at = [(Path("/roots/media/a.0123abcd.part"), _old(ORPHAN_AFTER_SECS))] + just_under = [(Path("/roots/media/a.0123abcd.part"), _old(ORPHAN_AFTER_SECS - 1))] + assert orphaned_parts(at, set(), NOW) == [Path("/roots/media/a.0123abcd.part")] assert orphaned_parts(just_under, set(), NOW) == [] @@ -243,9 +245,9 @@ def test_the_janitor_deletes_the_abandoned_and_keeps_the_rest(tmp_path): one somebody is still writing stay, and a finished file is never a candidate.""" root = _root(tmp_path, "media") - old = _aged(root.path / "abandoned.mkv.part", ORPHAN_AFTER_SECS + 60) - recent = _aged(root.path / "fresh.mkv.part", 30) - live = _aged(root.path / "sending.mkv.part", ORPHAN_AFTER_SECS * 2) + old = _aged(root.path / "abandoned.mkv.0123abcd.part", ORPHAN_AFTER_SECS + 60) + recent = _aged(root.path / "fresh.mkv.0123abcd.part", 30) + live = _aged(root.path / "sending.mkv.0123abcd.part", ORPHAN_AFTER_SECS * 2) finished = _aged(root.path / "done.mkv", ORPHAN_AFTER_SECS * 5) uploads = PartialUploads() @@ -262,7 +264,7 @@ def test_a_group_that_has_never_uploaded_anything_is_handled(tmp_path): """No `partial_uploads` in the context yet — it is created on first use, so a node that has been up for five minutes has none.""" root = _root(tmp_path, "media") - old = _aged(root.path / "left.mkv.part", ORPHAN_AFTER_SECS + 1) + old = _aged(root.path / "left.mkv.0123abcd.part", ORPHAN_AFTER_SECS + 1) daemon = _daemon({"g1": {"roots": RootSet(roots=[root])}}) assert daemon._reap_once() == 1 assert not old.exists() @@ -362,7 +364,7 @@ async def test_an_upload_in_flight_is_known_to_the_reaper(tmp_path): total_chunks=2)) live = ctx["partial_uploads"].live_paths() assert len(live) == 1 - assert next(iter(live)).name == "film.mkv.part" + assert re.fullmatch(r"film\.mkv\.[0-9a-f]{8}\.part", next(iter(live)).name) assert next(iter(live)).exists() @@ -489,3 +491,97 @@ async def test_an_upload_chunk_says_its_slot_is_in_use(tmp_path): assert _errors(peer) == [] assert slots.leases["up-1"].used is True, ( "the node still believes nobody took this slot up, and will reclaim it") + + + +# ── never replacing a file ────────────────────────────────────────────────── + +async def test_two_members_sending_one_name_at_once_get_two_files(tmp_path): + """ + Both chose a free name at chunk 0, and the free name was the same: the + final file did not exist yet, only the first one's `.part`. They then wrote + one `.part`, the second truncating the first, and published it twice. + """ + ctx = _group_ctx(tmp_path) + alice, bob = _peer(ctx, "alice"), _peer(ctx, "bob") + for who, data in ((alice, b"hers-1"), (bob, b"his-1")): + await who._do_file_upload(sealed_upload(who, filename="IMG_1234.jpg", data=data, + chunk_index=0, total_chunks=2)) + for who, data in ((alice, b"hers-2"), (bob, b"his-2")): + await who._do_file_upload(sealed_upload(who, filename="IMG_1234.jpg", data=data, + chunk_index=1, total_chunks=2)) + assert _errors(alice) == [] and _errors(bob) == [] + + root = ctx["roots"].roots[0].path + assert (root / "IMG_1234.jpg").read_bytes() == b"hers-1hers-2" + assert (root / "IMG_1234 (2).jpg").read_bytes() == b"his-1his-2" + assert _acks(bob, ctx)[-1]["stored_as"] == "IMG_1234 (2).jpg" + + +async def test_a_file_that_appears_during_an_upload_is_not_replaced(tmp_path): + """The operator copies a file in under the same name while a member's + upload is running. The rename at the end used to replace it.""" + ctx = _group_ctx(tmp_path) + peer = _peer(ctx) + await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"up-1", + chunk_index=0, total_chunks=2)) + root = ctx["roots"].roots[0].path + (root / "film.mkv").write_bytes(b"the operator's") + await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"up-2", + chunk_index=1, total_chunks=2)) + + assert _errors(peer) == [] + assert (root / "film.mkv").read_bytes() == b"the operator's" + assert (root / "film (2).mkv").read_bytes() == b"up-1up-2" + assert _acks(peer, ctx)[-1]["stored_as"] == "film (2).mkv" + assert not list(root.glob("*.part")), "the part was left behind" + + +def test_without_hard_links_a_file_is_still_not_replaced(tmp_path, monkeypatch): + """FAT, exFAT and some network shares have no hard links.""" + from meshbay_node import roots as roots_mod + + def no_links(*a, **k): + raise PermissionError("operation not permitted") + monkeypatch.setattr(roots_mod.os, "link", no_links) + (tmp_path / "a.txt").write_bytes(b"there first") + part = tmp_path / "a.txt.0123abcd.part" + part.write_bytes(b"uploaded") + + name = roots_mod.publish_upload(part, tmp_path, "a.txt", "a.txt") + assert name == "a (2).txt" + assert (tmp_path / "a.txt").read_bytes() == b"there first" + assert (tmp_path / "a (2).txt").read_bytes() == b"uploaded" + assert not part.exists() + + +# ── files Explorer acts on by itself ───────────────────────────────────────── + + +@pytest.mark.parametrize("name", ["desktop.ini", "Desktop.INI", "photos.lnk", "site.url", + "x.scf", "Docs.library-ms", "s.searchConnector-ms"]) +async def test_a_file_explorer_acts_on_is_refused(tmp_path, name): + ctx = _group_ctx(tmp_path) + peer = _peer(ctx) + await peer._do_file_upload(sealed_upload(peer, filename=name, data=b"[x]")) + assert [m.get("code") for m in _errors(peer)] == ["file_type_refused"] + assert not any(ctx["roots"].roots[0].path.iterdir()) + + +async def test_an_ordinary_file_with_a_near_name_is_accepted(tmp_path): + ctx = _group_ctx(tmp_path) + peer = _peer(ctx) + await peer._do_file_upload(sealed_upload(peer, filename="url-notes.txt", data=b"x")) + assert _errors(peer) == [] + + + +def test_a_part_the_node_did_not_write_is_never_deleted(): + """A browser's download in progress, a copy the operator is making: a + `.part` without the node's tag is somebody else's, however old.""" + doomed = orphaned_parts( + [(Path("/roots/media/film.mkv.part"), _old(ORPHAN_AFTER_SECS * 10)), + (Path("/roots/media/report.pdf.part"), _old(ORPHAN_AFTER_SECS * 10)), + (Path("/roots/media/film.mkv.0123abcd.part"), _old(ORPHAN_AFTER_SECS * 10))], + live=set(), now=NOW) + assert doomed == [Path("/roots/media/film.mkv.0123abcd.part")] diff --git a/packages/meshbay-node/tests/test_revocation_closes_sessions.py b/packages/meshbay-node/tests/test_revocation_closes_sessions.py new file mode 100644 index 0000000..7c8e79c --- /dev/null +++ b/packages/meshbay-node/tests/test_revocation_closes_sessions.py @@ -0,0 +1,53 @@ +""" +A revocation the hub signed closes what it revokes, not only what comes next. + +The denylist refuses the next connection. A revoked account's live sessions +were left open — streaming, downloading, chatting — until they happened to end; +only a group's revocation closed its sessions. +""" + +import asyncio + +import pytest +from meshbay_node.daemon import NodeDaemon +from meshbay_node.transport.quic_server import Denylist + + +class _Session: + def __init__(self, user_id, group_id): + self._user_id = user_id + self._group_id = group_id + self.closed = False + + async def close(self): + self.closed = True + + +def _daemon(sessions): + d = NodeDaemon.__new__(NodeDaemon) + d._webrtc = type("T", (), {"_sessions": sessions})() + return d + + +@pytest.mark.asyncio +async def test_a_revoked_account_is_disconnected(tmp_path): + mallory, alice = _Session("mallory", "g1"), _Session("alice", "g1") + phone = _Session("mallory", "g2") + d = _daemon({"a": mallory, "b": alice, "c": phone}) + deny = Denylist(path=tmp_path / "deny.json") + + d._apply_revocation(deny, "user", "mallory") + await asyncio.sleep(0.05) + + assert mallory.closed and phone.closed, "every session of the account ends" + assert not alice.closed + assert deny.is_denied("mallory", "") + + +@pytest.mark.asyncio +async def test_a_revoked_group_is_still_disconnected(tmp_path): + a, b = _Session("alice", "g1"), _Session("alice", "g2") + d = _daemon({"a": a, "b": b}) + d._apply_revocation(Denylist(path=tmp_path / "deny.json"), "group", "g1") + await asyncio.sleep(0.05) + assert a.closed and not b.closed |