aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/tests
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/tests')
-rw-r--r--packages/meshbay-node/tests/golden/cli.json9
-rw-r--r--packages/meshbay-node/tests/test_attach_from_the_hub.py99
-rw-r--r--packages/meshbay-node/tests/test_chat_is_bounded.py35
-rw-r--r--packages/meshbay-node/tests/test_ffprobe_is_bounded.py61
-rw-r--r--packages/meshbay-node/tests/test_linkpreview.py94
-rw-r--r--packages/meshbay-node/tests/test_member_capacity.py151
-rw-r--r--packages/meshbay-node/tests/test_member_errors_are_plain.py22
-rw-r--r--packages/meshbay-node/tests/test_ops.py2
-rw-r--r--packages/meshbay-node/tests/test_partial_uploads.py118
-rw-r--r--packages/meshbay-node/tests/test_revocation_closes_sessions.py53
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