diff options
Diffstat (limited to 'packages/meshbay-node')
32 files changed, 1148 insertions, 157 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/cli/groups.py b/packages/meshbay-node/src/meshbay_node/cli/groups.py index 3fab4b9..fca4800 100644 --- a/packages/meshbay-node/src/meshbay_node/cli/groups.py +++ b/packages/meshbay-node/src/meshbay_node/cli/groups.py @@ -64,7 +64,7 @@ def group(args) -> None: sys.exit(1) if not args.target or not args.dir: print("usage: meshbay-node group add <name> --dir <path> " - "[--no-writable]") + "[--no-writable] [--open]") print() print("The group must already exist on the hub and be yours. This") print("only tells the node to host it, and picks its first") @@ -78,12 +78,19 @@ def group(args) -> None: # is not a working group. Every root added *later* is read-only by # default, which is the opposite rule and the right one there. writable = args.writable is not False + # How people join is the operator's to say here, never read from the hub. + join_policy = "open" if getattr(args, "open", False) else "invite" body = {"name": args.target, "shared_dir": args.dir, - "writable": writable} + "writable": writable, "join_policy": join_policy} out = _daemon_api(cfg, "/api/groups/attach", method="POST", body=body) print(f"{out['name']} ({out['group_id'][:8]}) added to {out['config']}") print(f" shared_dir {out['shared_dir']}" f" ({'read-write' if writable else 'read-only'})") + print(f" join_policy {join_policy}") + if out.get("hub_join_policy") == "open" and join_policy != "open": + print() + print("The hub lists this group as open; this node admits by invitation") + print("only. To host it open, remove it and add it again with --open.") print() print("Tell the daemon to re-read its config, then give the group a key:") print(" meshbay-node reload") diff --git a/packages/meshbay-node/src/meshbay_node/cli/parser.py b/packages/meshbay-node/src/meshbay_node/cli/parser.py index 29a9f63..cfc4704 100644 --- a/packages/meshbay-node/src/meshbay_node/cli/parser.py +++ b/packages/meshbay-node/src/meshbay_node/cli/parser.py @@ -87,6 +87,9 @@ def build_parser() -> argparse.ArgumentParser: parser.add_argument("--no-removable", action="store_false", dest="removable", help="mark root as not removable (root set)") + parser.add_argument("--open", action="store_true", + help="group add: anyone the hub lists the group to may join " + "(default: by invitation only)") parser.add_argument("--name", default=None, help="root name (root add; defaults to directory basename)") parser.add_argument("--log-level", default="INFO", diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 45a6b9a..1777bc3 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -577,12 +577,12 @@ class NodeDaemon(EnrichmentMixin): log.info("QUIC server disabled ([node] quic_enabled = false)") # 8. Hub WebSocket (signaling + revocations + WebRTC offers) - async def on_webrtc_offer(sdp, peer_id, ice_candidates): + async def on_webrtc_offer(sdp, peer_id, ice_candidates, user_id=""): if not self._webrtc: return None try: answer_sdp, answer_ice = await self._webrtc.handle_offer( - sdp, peer_id) + sdp, peer_id, user_id) log.info("WebRTC answer for peer=%s (%d peers)", peer_id, self._webrtc.active_peers) return (answer_sdp, answer_ice) @@ -601,19 +601,8 @@ class NodeDaemon(EnrichmentMixin): payload = _jwt.decode( token, session.hub_pk_pem, algorithms=["EdDSA"], options={"verify_exp": False}) - target = payload.get("target") - tid = payload.get("target_id", "") - if target == "user": - denylist.deny_user(tid) - elif target == "group": - # H4: previously dropped on the floor, so "suspend a - # group" was a hub-only gesture that no node enforced. - denylist.deny_group(tid) - self._drop_group_sessions(tid) - elif target == "jti": - denylist.deny_jti(tid) - else: - log.warning("Unknown revocation target: %r", target) + self._apply_revocation(denylist, payload.get("target"), + payload.get("target_id", "")) except Exception as e: log.warning("Invalid revocation token: %s", e) @@ -1522,6 +1511,35 @@ class NodeDaemon(EnrichmentMixin): except Exception: pass + def _apply_revocation(self, denylist, target, target_id: str) -> None: + """ + What a revocation the hub signed does on this node. + + A revoked account or group is refused from now on, and its live sessions + are closed: a denylist entry alone stops the next connection and leaves + the current one streaming, downloading and chatting until it happens to + disconnect. + """ + if target == "user": + denylist.deny_user(target_id) + self._drop_user_sessions(target_id) + elif target == "group": + denylist.deny_group(target_id) + self._drop_group_sessions(target_id) + elif target == "jti": + denylist.deny_jti(target_id) + else: + log.warning("Unknown revocation target: %r", target) + + def _drop_user_sessions(self, user_id: str) -> None: + """Close every live session of a revoked account.""" + if not self._webrtc or not user_id: + return + for session in list(self._webrtc._sessions.values()): + if getattr(session, "_user_id", None) == user_id: + spawn(session.close()) + log.info("Dropped session for revoked account %s", user_id[:8]) + def _drop_group_sessions(self, group_id: str) -> None: """Close live sessions for a revoked group (H4).""" if not self._webrtc or not group_id: diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py index 2ce8993..54f87f0 100644 --- a/packages/meshbay-node/src/meshbay_node/hub_client.py +++ b/packages/meshbay-node/src/meshbay_node/hub_client.py @@ -451,7 +451,8 @@ class HubClient: """Negotiate one WebRTC offer and return the answer, off the read loop.""" try: answer = await on_webrtc_offer( - msg["sdp"], msg["peer_id"], msg.get("ice_candidates", [])) + msg["sdp"], msg["peer_id"], msg.get("ice_candidates", []), + str(msg.get("user_id") or "")) except Exception as e: log.warning("WebRTC offer from %s failed: %s", str(msg.get("peer_id"))[:8], e) diff --git a/packages/meshbay-node/src/meshbay_node/linkpreview.py b/packages/meshbay-node/src/meshbay_node/linkpreview.py index 6e1e618..493ac65 100644 --- a/packages/meshbay-node/src/meshbay_node/linkpreview.py +++ b/packages/meshbay-node/src/meshbay_node/linkpreview.py @@ -13,16 +13,19 @@ hub: apps (`_fetch_and_cache_poster`), over the same authorised path. Because the node makes an outbound request to an address a *member* chose, -this is an SSRF surface. `safe_url()` is the gate: http(s) only, no +this is an SSRF surface. `check_url()` is the gate: http(s) only, no credentials, the port restricted to the web set, and every resolved address must be globally routable — no loopback, private, link-local, multicast or -reserved range. Redirects are followed by hand so every hop is re-checked, -and the address the connection actually landed on is re-checked against the -same rule (`_reject_if_rebound`), so a name that resolves clean and then to -something internal (rebinding) does not get its body read. A full pin — -connect to the validated literal, verify the certificate for the name — is -the remaining hardening. How many previews a member can trigger is -rate-limited by the caller (`_do_link_preview_request`). +reserved range. Redirects are followed by hand so every hop is re-checked. + +**The connection goes to the address that was checked** (`_PinnedBackend`): +the name is resolved once, off the event loop, every answer is checked, and +the socket is opened to that IP literal — TLS still verifies the certificate +for the name. Resolving to check and letting the HTTP client resolve again to +connect would let a name answer clean the first time and with a LAN address +the second (DNS rebinding), and the request would be sent before anything +could look. How many previews a member can trigger is rate-limited by the +caller (`_do_link_preview_request`). Nothing is stored durably: the caller keeps an in-memory TTL cache and the OG image rides the existing `media_cache` thumb store (same as a poster). @@ -38,6 +41,7 @@ from html.parser import HTMLParser from io import BytesIO from urllib.parse import urljoin, urlsplit +import httpcore import httpx log = logging.getLogger(__name__) @@ -77,7 +81,13 @@ def _addr_is_public(ip: str) -> bool: def safe_url(url: str) -> str: - """Return the URL unchanged if it is safe to fetch, else raise UnsafeURL.""" + """ + The URL unchanged if its shape is safe to fetch, else raise UnsafeURL. + + Shape only — scheme, credentials, port, and the address when it is a + literal. A name is checked by `check_url` and, again, at connect time; + resolving here would block the event loop on a member's choice of name. + """ if not isinstance(url, str) or len(url) > 2048: raise UnsafeURL("missing or oversized") parts = urlsplit(url) @@ -94,23 +104,102 @@ def safe_url(url: str) -> str: raise UnsafeURL("bad port") if port is not None and port not in _ALLOWED_PORTS: raise UnsafeURL(f"port {port}") - # An IP literal is checked directly; a name is resolved and every answer - # must be public — a hostname with one public and one 127.0.0.1 record - # would otherwise be a way in. + if _is_literal(host) and not _addr_is_public(host): + raise UnsafeURL(f"non-public address {host}") + return url + + +def _is_literal(host: str) -> bool: + try: + ipaddress.ip_address(host) + return True + except ValueError: + return False + + +_RESOLVE_TIMEOUT = 5.0 + + +async def resolve_public(host: str, port: int) -> str: + """ + One public address for `host`, resolved off the event loop, or UnsafeURL. + + Every answer must be public: a name with one public and one 127.0.0.1 + record would otherwise be a way in. + """ + if _is_literal(host): + if not _addr_is_public(host): + raise UnsafeURL(f"non-public address {host}") + return host + loop = asyncio.get_running_loop() try: - infos = socket.getaddrinfo(host, parts.port or (443 if parts.scheme == "https" else 80), - proto=socket.IPPROTO_TCP) - except socket.gaierror as e: + infos = await asyncio.wait_for( + loop.getaddrinfo(host, port, proto=socket.IPPROTO_TCP), _RESOLVE_TIMEOUT) + except (socket.gaierror, TimeoutError) as e: raise UnsafeURL(f"cannot resolve: {e}") - resolved = {info[4][0] for info in infos} + resolved = list(dict.fromkeys(info[4][0] for info in infos)) if not resolved: raise UnsafeURL("resolves to nothing") bad = [ip for ip in resolved if not _addr_is_public(ip)] if bad: raise UnsafeURL(f"non-public address {bad[0]}") + return resolved[0] + + +async def check_url(url: str) -> str: + """`safe_url`, and the name's addresses checked too. The URL unchanged.""" + safe_url(url) + parts = urlsplit(url) + await resolve_public(parts.hostname or "", + parts.port or (443 if parts.scheme == "https" else 80)) return url +class _PinnedBackend(httpcore.AsyncNetworkBackend): + """ + Opens every connection to an address `resolve_public` checked. + + httpcore hands the backend the request's host; the TLS layer above still + uses that name for SNI and certificate verification, so pinning the socket + changes where it connects and nothing about whom it trusts. + """ + + def __init__(self) -> None: + self._inner = httpcore.AnyIOBackend() + + async def connect_tcp(self, host, port, timeout=None, local_address=None, + socket_options=None): + ip = await resolve_public(host, port) + return await self._inner.connect_tcp(ip, port, timeout=timeout, + local_address=local_address, + socket_options=socket_options) + + async def connect_unix_socket(self, path, timeout=None, socket_options=None): + raise UnsafeURL("no unix sockets") + + async def sleep(self, seconds: float) -> None: + await self._inner.sleep(seconds) + + +class _PinnedTransport(httpx.AsyncHTTPTransport): + """httpx's transport over a pool that connects through `_PinnedBackend`. + + No proxy from the environment (`trust_env=False`): a proxy would resolve + the name itself, and the pin would bind nothing. + """ + + def __init__(self) -> None: + super().__init__(trust_env=False, retries=0) + self._pool = httpcore.AsyncConnectionPool( + ssl_context=httpx.create_ssl_context(trust_env=False), + network_backend=_PinnedBackend(), max_connections=10) + + +def _new_client() -> httpx.AsyncClient: + return httpx.AsyncClient(transport=_PinnedTransport(), timeout=_TIMEOUT, + max_redirects=0, trust_env=False) + + class _HeadParser(HTMLParser): """Collects <title> text and name/property→content from <meta> in <head>. @@ -155,25 +244,6 @@ def _first(metas: dict[str, str], *keys: str) -> str | None: return None -def _reject_if_rebound(resp: httpx.Response) -> None: - """ - `safe_url` validated the name's addresses; this checks the one the - connection actually landed on, so a name that resolves clean and then to - something internal (DNS rebinding) does not get its body read. - - Best-effort: the `network_stream` extension is not present on every - transport (a MockTransport in tests has none), and its absence is not a - failure — the pre-check and the per-hop redirect re-check still stand. - """ - try: - stream = resp.extensions.get("network_stream") - addr = stream.get_extra_info("server_addr") if stream else None - except Exception: - return - if addr and not _addr_is_public(str(addr[0])): - raise UnsafeURL(f"connected to non-public address {addr[0]}") - - async def _get(client: httpx.AsyncClient, url: str) -> httpx.Response: """ One GET with manual, re-validated redirects, **body not read**. @@ -184,19 +254,14 @@ async def _get(client: httpx.AsyncClient, url: str) -> httpx.Response: compressed stream of any size was held in memory first. The caps are the only thing between a URL a member pasted and the node's memory. """ - current = safe_url(url) + current = await check_url(url) for _ in range(_MAX_REDIRECTS + 1): request = client.build_request("GET", current, headers={"User-Agent": _UA}) resp = await client.send(request, stream=True, follow_redirects=False) - try: - _reject_if_rebound(resp) - except Exception: - await resp.aclose() - raise if resp.is_redirect and "location" in resp.headers: location = resp.headers["location"] await resp.aclose() - current = safe_url(urljoin(current, location)) + current = await check_url(urljoin(current, location)) continue return resp raise UnsafeURL("too many redirects") @@ -234,7 +299,7 @@ async def fetch_preview(url: str, *, client: httpx.AsyncClient | None = None) -> """ own = client is None if own: - client = httpx.AsyncClient(timeout=_TIMEOUT, max_redirects=0) + client = _new_client() try: return await asyncio.wait_for(_preview(client, url), _TOTAL_DEADLINE) except (httpx.HTTPError, UnsafeURL, TimeoutError) as e: @@ -246,7 +311,6 @@ async def fetch_preview(url: str, *, client: httpx.AsyncClient | None = None) -> async def _preview(client: httpx.AsyncClient, url: str) -> dict | None: - safe_url(url) resp = await _get(client, url) try: ctype = resp.headers.get("content-type", "").split(";")[0].strip().lower() @@ -272,7 +336,7 @@ async def _preview(client: httpx.AsyncClient, url: str) -> dict | None: if image: image = urljoin(final_url, image) try: - safe_url(image) + await check_url(image) except UnsafeURL: image = None @@ -294,7 +358,7 @@ async def fetch_image(url: str, *, client: httpx.AsyncClient | None = None) -> b """Fetch and re-encode an OG image to a small JPEG. None on any failure.""" own = client is None if own: - client = httpx.AsyncClient(timeout=_TIMEOUT, max_redirects=0) + client = _new_client() try: raw = await asyncio.wait_for(_image_bytes(client, url), _TOTAL_DEADLINE) except (httpx.HTTPError, UnsafeURL, TimeoutError) as e: @@ -310,7 +374,6 @@ async def fetch_image(url: str, *, client: httpx.AsyncClient | None = None) -> b async def _image_bytes(client: httpx.AsyncClient, url: str) -> bytes | None: - safe_url(url) resp = await _get(client, url) try: ctype = resp.headers.get("content-type", "").split(";")[0].strip().lower() diff --git a/packages/meshbay-node/src/meshbay_node/media_probe.py b/packages/meshbay-node/src/meshbay_node/media_probe.py index a6267d9..e9090ea 100644 --- a/packages/meshbay-node/src/meshbay_node/media_probe.py +++ b/packages/meshbay-node/src/meshbay_node/media_probe.py @@ -10,6 +10,12 @@ import asyncio import json from dataclasses import dataclass, field +# How long ffprobe may take over one file's headers. The file is a member's +# upload as often as the operator's own: one that keeps ffprobe busy must not +# keep the stream request, the subtitle request or the enrichment slot that +# asked for it — the same bound the index-time enrichment already put around it. +FFPROBE_TIMEOUT_SECS = 30 + _H264_PROFILES = {"Baseline": "42", "Main": "4d", "High": "64", "High 10": "6e"} # Source video codecs whose MSE codec string is real but which no mainstream @@ -158,7 +164,16 @@ async def probe_video(path: str) -> VideoProbe: "-of", "json", path, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, ) - stdout, _ = await proc.communicate() + try: + stdout, _ = await asyncio.wait_for(proc.communicate(), FFPROBE_TIMEOUT_SECS) + except (TimeoutError, asyncio.CancelledError) as e: + # Killed, not abandoned: a cancelled wait leaves the process running, + # and a caller's own timeout (enrich.py) cancels exactly this wait. + proc.kill() + await proc.wait() + if isinstance(e, asyncio.CancelledError): + raise + raise RuntimeError(f"ffprobe timed out after {FFPROBE_TIMEOUT_SECS}s") from None info = json.loads(stdout) duration = float(info.get("format", {}).get("duration", 0)) diff --git a/packages/meshbay-node/src/meshbay_node/ops/groups.py b/packages/meshbay-node/src/meshbay_node/ops/groups.py index 1d8003c..ac14c3a 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/groups.py +++ b/packages/meshbay-node/src/meshbay_node/ops/groups.py @@ -9,7 +9,7 @@ from meshbay_common.crypto import generate_gek, wrap_gek_aes from meshbay_node.config import DEFAULT_CONFIG_PATH from meshbay_node.ops.core import OpError, _config, _group_ctx, _hub -from meshbay_node.ops.node_toml import _find_group_range +from meshbay_node.ops.node_toml import _find_group_range, toml_string log = logging.getLogger("meshbay_node.ops") @@ -142,17 +142,28 @@ async def list_groups(state: dict) -> dict: return {"groups": out, "operator_paired": has_operator, "settings": settings} +JOIN_POLICIES = ("invite", "open") + + async def attach_group(state: dict, name: str, shared_dir: str, - writable: bool = True) -> dict: + writable: bool = True, join_policy: str = "invite") -> dict: """ Write a new [[groups]] block into node.toml. The name-to-id lookup happens here because this process is the one logged into the hub. Nothing is created on the hub: the group already exists, this only tells the node to host it. + + `join_policy` is the operator's, given with this request, and `invite` + unless they say otherwise. The hub's own record of the group is not read + for it: a hub that could declare a group open would be handed its key by + anyone it sent. The hub's value is returned beside it, so a caller can say + when the two differ. """ if not name or not shared_dir: raise OpError("name and shared_dir are required") + if join_policy not in JOIN_POLICIES: + raise OpError(f"join_policy must be one of {', '.join(JOIN_POLICIES)}") config = _config(state) hub = _hub(state) try: @@ -182,12 +193,12 @@ async def attach_group(state: dict, name: str, shared_dir: str, raise OpError(f"Cannot create {path}: {e}") from e conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) - join_policy = group.get("join_policy", "invite") + visibility = "public" if join_policy == "open" else "private" block = (f'\n[[groups]]\n' - f'id = "{group["id"]}"\n' - f'name = "{group["name"]}"\n' - f'visibility = "{group.get("visibility", "private")}"\n' - f'join_policy = "{join_policy}"\n') + f'id = {toml_string(group["id"])}\n' + f'name = {toml_string(group["name"])}\n' + f'visibility = {toml_string(visibility)}\n' + f'join_policy = {toml_string(join_policy)}\n') # No `upload_dir` here. `GroupConfig.__post_init__` still *reads* it, so an # existing node.toml keeps working — but what it does on read is force every # other root read-only and append that path as the one writable one, which @@ -198,7 +209,7 @@ async def attach_group(state: dict, name: str, shared_dir: str, block += (f'\n [[groups.roots]]\n' # Forward slashes: a Windows path in a TOML basic string is a # parse error (`\U`, `\a`, ... are escapes). pathlib reads `/`. - f' path = "{path.as_posix()}"\n' + f' path = {toml_string(path.as_posix())}\n' f' writable = {"true" if writable else "false"}\n') try: with conf_path.open("a", encoding="utf-8", newline="\n") as f: @@ -208,7 +219,8 @@ async def attach_group(state: dict, name: str, shared_dir: str, result = {"group_id": group["id"], "name": group["name"], "shared_dir": str(path), "config": str(conf_path), - "writable": writable, + "writable": writable, "join_policy": join_policy, + "hub_join_policy": group.get("join_policy", "invite"), "note": "restart the node to pick it up"} return result diff --git a/packages/meshbay-node/src/meshbay_node/ops/node_toml.py b/packages/meshbay-node/src/meshbay_node/ops/node_toml.py index f711f26..2407722 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/node_toml.py +++ b/packages/meshbay-node/src/meshbay_node/ops/node_toml.py @@ -3,14 +3,50 @@ from __future__ import annotations import re +import tomllib from pathlib import Path from meshbay_node.ops.core import OpError +def toml_string(value: str) -> str: + """A TOML basic string holding `value` exactly, quotes included. + + Every string written into node.toml goes through here. A value with a quote + or a newline in it — a group name, a folder name, any of them chosen by + someone else — would otherwise end the string and write lines of its own. + """ + out = ['"'] + for ch in str(value): + if ch == '"': + out.append('\\"') + elif ch == "\\": + out.append("\\\\") + elif ord(ch) < 0x20 or ord(ch) == 0x7F: + out.append(f"\\u{ord(ch):04x}") + else: + out.append(ch) + out.append('"') + return "".join(out) + + +def _string_value(line: str, key: str) -> str | None: + """The string `key` holds on this line, unescaped — or None. + + Read as TOML, not by pattern: a value written by `toml_string` may carry an + escaped quote or backslash, which a `"([^"]*)"` pattern would cut short. + """ + if not re.match(r"^\s*" + re.escape(key) + r"\s*=", line): + return None + try: + value = tomllib.loads(line.strip()).get(key) + except tomllib.TOMLDecodeError: + return None + return value if isinstance(value, str) else None + + def _find_group_range(lines: list[str], group_id: str) -> tuple[int, int] | None: """Line range of a [[groups]] block by id: (start, end_exclusive).""" - id_re = re.compile(r'^\s*id\s*=\s*"([^"]*)"') block_starts: list[int] = [] for i, line in enumerate(lines): if line.strip() == "[[groups]]": @@ -24,8 +60,7 @@ def _find_group_range(lines: list[str], group_id: str) -> tuple[int, int] | None boundary = k break for k in range(start + 1, boundary): - m = id_re.match(lines[k]) - if m and m.group(1) == group_id: + if _string_value(lines[k], "id") == group_id: return (start, boundary) return None @@ -61,8 +96,10 @@ def _update_node_toml(conf_path: Path, updates: dict) -> None: if isinstance(value, bool): return f"{key} = {'true' if value else 'false'}" if isinstance(value, list): - items = ", ".join(f'"{v}"' for v in value) + items = ", ".join(toml_string(v) for v in value) return f"{key} = [{items}]" + if isinstance(value, str): + return f"{key} = {toml_string(value)}" return f"{key} = {value}" remaining = dict(updates) @@ -115,7 +152,6 @@ def _remove_roots_block(conf_path: Path, group_id: str, raise OpError(f"Group {group_id[:8]} not found in {conf_path}") start, end = rng - path_re = re.compile(r'^\s*path\s*=\s*"([^"]*)"') roots_starts: list[int] = [] for i in range(start + 1, end): if lines[i].strip() == "[[groups.roots]]": @@ -124,10 +160,10 @@ def _remove_roots_block(conf_path: Path, group_id: str, for j, rs in enumerate(roots_starts): rs_end = roots_starts[j + 1] if j + 1 < len(roots_starts) else end for k in range(rs, rs_end): - m = path_re.match(lines[k]) - if m: + raw = _string_value(lines[k], "path") + if raw is not None: try: - p = str(Path(m.group(1)).expanduser().resolve()) + p = str(Path(raw).expanduser().resolve()) except OSError: continue if p == resolved_path: @@ -153,7 +189,6 @@ def _update_root_field(conf_path: Path, group_id: str, raise OpError(f"Group {group_id[:8]} not found in {conf_path}") start, end = rng - path_re = re.compile(r'^\s*path\s*=\s*"([^"]*)"') writable_re = re.compile(r'^\s*(writable|upload)\s*=') removable_re = re.compile(r'^\s*removable\s*=') roots_starts: list[int] = [] @@ -165,10 +200,10 @@ def _update_root_field(conf_path: Path, group_id: str, rs_end = roots_starts[j + 1] if j + 1 < len(roots_starts) else end found_path = False for k in range(rs, rs_end): - m = path_re.match(lines[k]) - if m: + raw = _string_value(lines[k], "path") + if raw is not None: try: - p = str(Path(m.group(1)).expanduser().resolve()) + p = str(Path(raw).expanduser().resolve()) except OSError: continue if p == resolved_path: diff --git a/packages/meshbay-node/src/meshbay_node/ops/roots.py b/packages/meshbay-node/src/meshbay_node/ops/roots.py index e3e2781..480fe63 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/roots.py +++ b/packages/meshbay-node/src/meshbay_node/ops/roots.py @@ -8,7 +8,12 @@ from pathlib import Path from meshbay_node.config import DEFAULT_CONFIG_PATH from meshbay_node.ops.core import OpError, _config, _group_ctx, _roster -from meshbay_node.ops.node_toml import _insert_roots_block, _remove_roots_block, _update_root_field +from meshbay_node.ops.node_toml import ( + _insert_roots_block, + _remove_roots_block, + _update_root_field, + toml_string, +) from meshbay_node.roots import RootError, RootSet, off_disk log = logging.getLogger("meshbay_node.ops") @@ -46,11 +51,11 @@ async def add_root(state: dict, group_id: str, path: str, *, raise OpError(f"Cannot create {added.path}: {e}") from e conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) - root_block = f' [[groups.roots]]\n path = "{added.path.as_posix()}"' + root_block = f' [[groups.roots]]\n path = {toml_string(added.path.as_posix())}' if name: - root_block += f'\n name = "{added.name}"' + root_block += f'\n name = {toml_string(added.name)}' if kind != "generic": - root_block += f'\n kind = "{added.kind}"' + root_block += f'\n kind = {toml_string(added.kind)}' if writable: root_block += '\n writable = true' if removable: diff --git a/packages/meshbay-node/src/meshbay_node/roots.py b/packages/meshbay-node/src/meshbay_node/roots.py index 89b0441..2de0708 100644 --- a/packages/meshbay-node/src/meshbay_node/roots.py +++ b/packages/meshbay-node/src/meshbay_node/roots.py @@ -30,6 +30,7 @@ from __future__ import annotations import asyncio import logging +import os import re from concurrent.futures import ThreadPoolExecutor from dataclasses import dataclass, field @@ -52,25 +53,77 @@ SAFE_UPLOAD_NAME = re.compile( re.UNICODE) -def _free_name(directory: Path, filename: str) -> str: +# Files Windows Explorer acts on by itself when it shows a folder: a link's +# icon, a folder's settings, a search connector. Placed by a member in a folder +# the operator browses, any of them can make Explorer contact a server of the +# member's choosing with the operator's Windows credentials — a known attack, +# and why mail providers refuse the same types. Refused for every node: a +# Linux node's folder may be shared to Windows machines. +SHELL_ACTIVE_NAMES = frozenset({"desktop.ini"}) +SHELL_ACTIVE_SUFFIXES = (".lnk", ".url", ".scf", ".library-ms", ".searchconnector-ms") + + +def shell_active(filename: str) -> bool: + name = filename.lower() + return name in SHELL_ACTIVE_NAMES or name.endswith(SHELL_ACTIVE_SUFFIXES) + + +def _free_name(directory: Path, filename: str, + taken: frozenset[str] | set[str] = frozenset()) -> str: """ `filename`, or the first "name (n).ext" that is not taken. - Never returns the name of a file that exists, so an upload cannot replace - one — the property the per-user quarantine used to provide (C5a). + Never returns the name of a file that exists, nor one in `taken` — names + uploads in flight will publish under — so an upload cannot replace a file + or another upload (C5a). """ - if not (directory / filename).exists(): + def free(name: str) -> bool: + return name not in taken and not (directory / name).exists() + + if free(filename): return filename stem, dot, ext = filename.rpartition(".") if not dot: stem, ext = filename, "" for n in range(2, 1000): candidate = f"{stem} ({n}){dot}{ext}" - if not (directory / candidate).exists(): + if free(candidate): return candidate raise FileExistsError(filename) +def publish_upload(part: Path, directory: Path, stored_name: str, filename: str, + taken: frozenset[str] | set[str] = frozenset()) -> str: + """ + Move a finished `.part` to its name without ever replacing a file. The name + it was published under, which may not be `stored_name`. + + A rename replaces whatever is at the target, and the target can appear + while the upload runs — the operator copying a file in, another group's + upload into a shared folder. A hard link refuses an existing target, so it + is the publication; where the filesystem has none (FAT, exFAT, some network + shares), the existence check and the rename are as close as it gets. A + taken name moves on to the next free one rather than failing the upload. + """ + name = stored_name + for _ in range(8): + target = directory / name + try: + os.link(part, target) + except FileExistsError: + name = _free_name(directory, filename, taken) + continue + except OSError: + if target.exists(): + name = _free_name(directory, filename, taken) + continue + part.rename(target) + return name + part.unlink() + return name + raise FileExistsError(stored_name) + + def safe_subdir(roots: RootSet, rel: str) -> Path | None: """ Resolve a client-supplied directory inside one of the group's roots, or refuse. diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py index 857db33..a2d9f47 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py @@ -120,7 +120,7 @@ class MusicMixin: blob = await _transcode_audio_to_aac(file_path) except Exception as e: log.warning("Audio transcode failed for %s: %s", entry.id[:12], e) - self._send({"type": "error", "detail": f"Transcode failed: {e}"}) + self._send({"type": "error", "detail": "This track could not be converted"}) return transcode_hash = blake3.blake3(blob).hexdigest() diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py index 4337e24..d6a5248 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py @@ -2,6 +2,7 @@ stream a session holds, and the ffmpeg pipeline behind it.""" import asyncio +import contextlib import logging import time @@ -34,6 +35,10 @@ log = logging.getLogger("meshbay_node.transport.webrtc_server") # so the operator sets `max_concurrent_streams` under [node] in node.toml. This # value applies when they have said nothing. MAX_CONCURRENT_TRANSCODES = 8 +# Subtitle extractions one account may run at once. The player asks for one +# track at a time; two covers a quick change of track. Each holds a transcode +# slot for up to fifteen minutes on a long film. +MAX_SUBTITLE_JOBS_PER_ACCOUNT = 2 STREAM_SEGMENT_SIZE = 256 * 1024 @@ -179,6 +184,36 @@ class StreamingMixin: self._stream_task = asyncio.current_task() await self._stream_video(msg) + @contextlib.contextmanager + def _account_share(self, kind: str, limit: int): + """ + Hold one of this account's `limit` places for `kind`, or yield False. + + Counted on the node, across every session of the account: a member's + devices and tabs share one allowance. The node's own account is not + counted — it is the operator's machine. + """ + user = getattr(self, "_user_id", "") or "" + if not user or user == self._ctx.get("node_user_id"): + yield True + return + held = self._ctx.setdefault(f"_{kind}_by_account", {}) + if held.get(user, 0) >= limit: + yield False + return + held[user] = held.get(user, 0) + 1 + try: + yield True + finally: + held[user] -= 1 + if held[user] <= 0: + held.pop(user, None) + + def _streams_per_account(self) -> int: + """Half the node's viewers, rounded up: three screens in one home fit, + and no member alone takes every slot the operator set.""" + return max(1, -(-self._stream_capacity() // 2)) + def _transcode_semaphore(self) -> asyncio.Semaphore: """The node's stream budget, shared across every peer. @@ -205,22 +240,28 @@ class StreamingMixin: ctx = self._ctx log.info("stream: waiting for a slot (%d of %d in use)", ctx.get("_streams_in_flight", 0), self._stream_capacity()) - async with sem: - # Counted here rather than read back out of the semaphore's private - # `_value`: `set_capacity` needs to know how many slots are held in - # order to resize without letting the pool overshoot, and a number - # this code maintains itself is one that survives the semaphore - # object being replaced underneath it. - ctx["_streams_in_flight"] = ctx.get("_streams_in_flight", 0) + 1 - log.info("stream: slot acquired (%d of %d in use)", - ctx["_streams_in_flight"], self._stream_capacity()) - try: - await self._stream_video_inner(msg) - finally: - ctx["_streams_in_flight"] = max( - 0, ctx.get("_streams_in_flight", 1) - 1) - log.info("stream: slot released (%d of %d in use)", + with self._account_share("streams", self._streams_per_account()) as ok: + if not ok: + self._send({"type": "error", + "detail": "Too many videos playing from this account, " + "stop one and retry"}) + return + async with sem: + # Counted here rather than read back out of the semaphore's private + # `_value`: `set_capacity` needs to know how many slots are held in + # order to resize without letting the pool overshoot, and a number + # this code maintains itself is one that survives the semaphore + # object being replaced underneath it. + ctx["_streams_in_flight"] = ctx.get("_streams_in_flight", 0) + 1 + log.info("stream: slot acquired (%d of %d in use)", ctx["_streams_in_flight"], self._stream_capacity()) + try: + await self._stream_video_inner(msg) + finally: + ctx["_streams_in_flight"] = max( + 0, ctx.get("_streams_in_flight", 1) - 1) + log.info("stream: slot released (%d of %d in use)", + ctx["_streams_in_flight"], self._stream_capacity()) def _stream_capacity(self) -> int: return self._ctx.get("max_concurrent_streams") or MAX_CONCURRENT_TRANSCODES @@ -246,7 +287,10 @@ class StreamingMixin: try: probe = await _probe_video(str(file_path)) except Exception as e: - self._send({"type": "error", "detail": f"Probe failed: {e}"}) + # The cause to the operator's log; to the member, that it failed. + # ffmpeg's own words carry the operator's paths and versions. + log.warning("stream: probe failed for %s: %s", entry.id[:12], e) + self._send({"type": "error", "detail": "This video could not be read"}) return codec_str = probe.codec duration = probe.duration diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py index 70781eb..525f2a9 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py @@ -10,6 +10,7 @@ from meshbay_common.protocol import MNP from meshbay_node.media_probe import probe_video as _probe_video from meshbay_node.roots import off_disk +from meshbay_node.transport.webrtc.apps.streaming import MAX_SUBTITLE_JOBS_PER_ACCOUNT from meshbay_node.transport.webrtc.disk import _locate from meshbay_node.transport.webrtc.media_tools import ( _extract_subtitle_to_webvtt, @@ -135,10 +136,16 @@ class SubtitlesMixin: return budget = _subtitle_timeout_for(entry.size) - async with sem: - log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)", - file_id[:12], ordinal, budget) - blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget) + with self._account_share("subtitles", MAX_SUBTITLE_JOBS_PER_ACCOUNT) as ok: + if not ok: + log.info("subtitle: refused, account at its extraction share") + self._send({"type": "error", + "detail": "Subtitles are already being prepared, retry shortly"}) + return + async with sem: + log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)", + file_id[:12], ordinal, budget) + blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget) subtitle_hash = blake3.blake3(blob).hexdigest() await media_cache.put_thumb(subtitle_hash, synthetic_id, blob) @@ -158,8 +165,10 @@ class SubtitlesMixin: except BaseException as e: log.warning("subtitle: extract failed file=%s track=%d after %.1fs: %r", file_id[:12], ordinal, time.monotonic() - t0, e) + # The cause is in the log line above; ffmpeg's own words carry the + # operator's paths and versions. self._send({"type": "error", - "detail": f"Subtitle extraction failed: {e}"}) + "detail": "These subtitles could not be extracted"}) if isinstance(e, asyncio.CancelledError): raise finally: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py index 107a43e..31cc83c 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py @@ -7,7 +7,7 @@ import struct import msgpack from aiortc import RTCPeerConnection -from meshbay_node.transport.webrtc.limits import MAX_MSG +from meshbay_node.transport.webrtc.limits import MAX_MSG, UNPACK_LIMITS def _extract_dtls_fingerprint(sdp: str) -> bytes: @@ -77,7 +77,7 @@ class _DataChannelBuffer: break msg_bytes = bytes(self._buf[4:4 + length]) del self._buf[:4 + length] - yield msgpack.unpackb(msg_bytes, raw=False) + yield msgpack.unpackb(msg_bytes, raw=False, **UNPACK_LIMITS) def _get_remote_ip(pc: RTCPeerConnection) -> str: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py index 26ec27c..493d7f0 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py @@ -283,7 +283,18 @@ class ChatMixin: # anyone on the node. gctx = self._group_ctx() chat_store = gctx.get("chat_store") + # The two fields that travel in clear beside the ciphertext (the sealed + # envelope carries its own). Stored and relayed to every member, so they + # are what they claim to be and no larger: a name as long as a username, + # a thread id as long as a message id. Anything else is dropped. sender_name = msg.get("sender_name", "") + if not isinstance(sender_name, str) or len(sender_name) > 64: + sender_name = "" + thread_id = msg.get("thread_id") + id_like = (isinstance(thread_id, int) and not isinstance(thread_id, bool) + or isinstance(thread_id, str) and len(thread_id) <= 64) + if thread_id is not None and not id_like: + thread_id = None # Two shapes, and keeping them apart is what makes this deployable. # @@ -329,7 +340,7 @@ class ChatMixin: self._spawn(self._store_chat_message( chat_store, iteration=msg.get("iteration", 0), payload=raw, - thread_id=msg.get("thread_id"), sender_name=sender_name, + thread_id=thread_id, sender_name=sender_name, format=fmt, epoch=epoch, device=device, nonce=nonce, sig=sig, )) @@ -340,7 +351,7 @@ class ChatMixin: "sender_id": self._user_id, "sender_name": sender_name, "payload": payload, - "thread_id": msg.get("thread_id"), + "thread_id": thread_id, "timestamp": time.time(), "format": fmt, "epoch": epoch, @@ -598,7 +609,7 @@ class ChatMixin: the client asks, the node produces on demand, the asking device caches — nothing durable here). - `linkpreview.safe_url` is the SSRF gate: the URL a *member* chose + `linkpreview.check_url` is the SSRF gate: the URL a *member* chose decides an outbound request from the operator's machine, so http(s) only and the resolved address must be globally routable. Failure of any kind — blocked, unreachable, not HTML, nothing worth showing — diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py index 882ddbc..ae3f0e1 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py @@ -139,7 +139,21 @@ class SessionCore: log.info("WebRTC data received: %d bytes, msg #%d (peer=%s)", len(message), self._msg_count, self._peer_id) self._buffer.feed(message) - for msg in self._buffer.messages(): + decoded = self._buffer.messages() + while True: + try: + msg = next(decoded) + except StopIteration: + break + except ValueError as e: + # Over the size limit, or a container past its decode + # limit. The buffer still starts with that frame, so every + # later message would fail the same way: the session ends + # here. Only decoding is caught — a handler's own error is + # not a reason to drop the peer. + log.warning("Closing peer %s: %s", self._peer_id, e) + self._spawn(self.close()) + break self._handle_message(msg) if _WEBRTC_TRACE: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py index 7d458f4..86412e0 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py @@ -2,7 +2,21 @@ CHUNK_SIZE = 1024 * 1024 -MAX_MSG = 64 * 1024 * 1024 +# The largest message a peer may send once it has proved the group key. The +# largest a client really sends is a sealed playlist blob, 1 MiB (blobs.py); +# chat is 64 KiB and an upload chunk 48 KiB. Eight times the largest, because a +# message of many small objects decodes to several times its size in memory. +MAX_MSG = 8 * 1024 * 1024 + +# Per container, when a message is decoded: nothing a client sends comes near +# them, and without them one message of tiny elements is one enormous list. +UNPACK_LIMITS = { + "max_array_len": 100_000, + "max_map_len": 10_000, + "max_str_len": 1024 * 1024, + "max_bin_len": MAX_MSG, + "max_ext_len": 0, +} # What the `tr` on a chunk request turned out to be (see `_lease_of`). diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py index 7ecb0a5..22c1690 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py @@ -165,6 +165,7 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa fd, tmp_name = tempfile.mkstemp(suffix=".mp4") os.close(fd) tmp_path = Path(tmp_name) + proc = probe = None try: proc = await asyncio.create_subprocess_exec( platform.ffmpeg_cmd(), "-hide_banner", "-loglevel", "error", "-y", @@ -189,6 +190,12 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa log.warning("stream: seek probe failed at %.1fs: %r", t, e) return None finally: + # A timed-out wait leaves its process running; it is stopped here, not + # left to finish a seek nobody is waiting for. + for p in (proc, probe): + if p is not None and p.returncode is None: + p.kill() + await p.wait() await _discard_scratch(tmp_path) text = stdout.decode(errors="replace").strip().rstrip(",") try: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py index 1c6d1ce..02a50af 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py @@ -8,7 +8,14 @@ from pathlib import Path from meshbay_common.protocol import UPLOAD_PROBE_INDEX, file_upload_ack_wire, file_upload_payload from meshbay_node import uploads as uploads_mod -from meshbay_node.roots import SAFE_UPLOAD_NAME, RootSet, _free_name, off_disk +from meshbay_node.roots import ( + SAFE_UPLOAD_NAME, + RootSet, + _free_name, + off_disk, + publish_upload, + shell_active, +) from meshbay_node.transport.webrtc.disk import _append_chunk from meshbay_node.transport.webrtc.limits import LEASE_NONE, LEASE_QUEUED @@ -212,6 +219,9 @@ class UploadMixin: if not SAFE_UPLOAD_NAME.match(filename): _refuse("Invalid filename", "invalid_filename") return + if shell_active(filename): + _refuse("This type of file is not accepted", "file_type_refused") + return roots: RootSet | None = ctx.get("roots") if not roots: @@ -306,9 +316,13 @@ class UploadMixin: # A shared directory means two people can send the same name. Refusing the # second is safe but silly — everyone's camera produces IMG_1234.jpg — so # a free name is found instead. Never a replacement. + # Names other uploads into this directory will publish under are taken + # too: none of them is on disk yet. + reserved = uploads.reserved_names(rel_dir) stored_name = (state.stored_name if state - else await off_disk(roots, _free_name, target_dir, filename)) - tmp_path = target_dir / f"{stored_name}{uploads_mod.PART_SUFFIX}" + else await off_disk(roots, _free_name, target_dir, filename, reserved)) + tmp_path = (state.part_path if state and state.part_path + else target_dir / uploads_mod.part_name(stored_name)) final_path = target_dir / stored_name if chunk_index == UPLOAD_PROBE_INDEX: @@ -361,6 +375,22 @@ class UploadMixin: await off_disk(roots, _append_chunk, tmp_path, chunk_bytes, chunk_index == 0) uploads.advance(user_id, rel_dir, filename, chunk_index, len(chunk_bytes)) + last = chunk_index + 1 >= total_chunks + if last: + # Published before the last ack, so the ack names the file as it is + # on disk: publication never replaces a file, and may have had to + # take another free name for this one. + uploads.drop(user_id, rel_dir, filename) + try: + stored_name = await off_disk(roots, publish_upload, tmp_path, target_dir, + stored_name, filename, + uploads.reserved_names(rel_dir)) + except OSError as e: + log.warning("Upload %s could not be published: %s", stored_name, e) + _refuse("The file could not be stored", "store_failed") + return + final_path = target_dir / stored_name + self._send(file_upload_ack_wire( gek, self._group_id or "", upload_id=upload_id, @@ -372,9 +402,7 @@ class UploadMixin: dir=rel_dir, )) - if chunk_index + 1 >= total_chunks: - uploads.drop(user_id, rel_dir, filename) - await off_disk(roots, tmp_path.rename, final_path) + if last: log.info("Upload complete: %s (%d chunks, %d bytes)", stored_name, total_chunks, state.bytes) self._audit("file_upload", f"{rel_dir}/{stored_name}") diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py index 20636fa..0ca4b3d 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -70,7 +70,15 @@ log = logging.getLogger(__name__) # so the cost to an operator grew with the number of people in their groups. # Sized to be unreachable in ordinary use: a browser holds one connection per # open group, and a handshake unfinished after a minute is not going to finish. -MAX_PEER_SESSIONS = 64 +MAX_PEER_SESSIONS = 128 +# One account's share of them: half. A member of twenty groups hosted here, +# with three devices and a spare tab, holds up to 52 (each device keeps up to +# twelve connections for search and music, plus the open group page), so the +# share never refuses real use — and no single member can hold more than half +# of what the node will take. The node's own account is not counted against it: +# it is the operator's machine. Measured, an idle connected session costs about +# 0.15 MiB and one file descriptor. +MAX_PEER_SESSIONS_PER_ACCOUNT = MAX_PEER_SESSIONS // 2 UNAUTHENTICATED_SESSION_TIMEOUT = 60 # seconds @@ -238,8 +246,13 @@ class WebRTCTransport: pass return + def _sessions_of(self, user_id: str) -> int: + return sum(1 for s in self._sessions.values() + if (getattr(s, "_offer_user", "") or getattr(s, "_user_id", "") or "") + == user_id) + async def handle_offer( - self, offer_sdp: str, peer_id: str, + self, offer_sdp: str, peer_id: str, user_id: str = "", ) -> tuple[str, list[dict]]: """ Process a WebRTC SDP offer from a browser client. @@ -267,9 +280,18 @@ class WebRTCTransport: log.warning("Refusing WebRTC offer: %d peer sessions already open", len(self._sessions)) raise RuntimeError("Node is at its peer-connection limit") + # The account the hub authenticated for this offer. A hub that lied + # could only move the count between accounts; it can refuse offers + # outright already. + if (user_id and user_id != self._ctx.get("node_user_id") + and self._sessions_of(user_id) >= MAX_PEER_SESSIONS_PER_ACCOUNT): + log.warning("Refusing WebRTC offer: account %s already holds %d sessions", + user_id[:8], MAX_PEER_SESSIONS_PER_ACCOUNT) + raise RuntimeError("Account is at its peer-connection share") pc = RTCPeerConnection(configuration=config) session = WebRTCPeerSession(pc, self._ctx, peer_id=peer_id) + session._offer_user = user_id self._sessions[peer_id] = session self._reap_if_unauthenticated(peer_id) diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index bbc4649..f810903 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -179,6 +179,7 @@ def create_ui_app(state: dict) -> FastAPI: (payload.get("name") or "").strip(), (payload.get("shared_dir") or "").strip(), writable=bool(payload.get("writable", True)), + join_policy=str(payload.get("join_policy") or "invite"), )) reload_fn = state.get("reload_fn") if reload_fn: diff --git a/packages/meshbay-node/src/meshbay_node/uploads.py b/packages/meshbay-node/src/meshbay_node/uploads.py index dd31e1e..c485e39 100644 --- a/packages/meshbay-node/src/meshbay_node/uploads.py +++ b/packages/meshbay-node/src/meshbay_node/uploads.py @@ -24,6 +24,8 @@ build first. from __future__ import annotations +import re +import secrets import time from collections.abc import Iterable from dataclasses import dataclass, field @@ -34,6 +36,26 @@ from pathlib import Path # recognise one, and a second spelling of it would be a bug nobody could see. PART_SUFFIX = ".part" + +def part_name(stored_name: str) -> str: + """The `.part` one upload writes: its final name, a tag of its own, `.part`. + + Its own, because the final name alone is shared: two uploads that settled + on one name — two groups hosting one folder, each with its own lock — + would write one file, the second truncating the first. + """ + return f"{stored_name}.{secrets.token_hex(4)}{PART_SUFFIX}" + + +# What `part_name` writes, and the only thing the reaper deletes. A `.part` +# without the node's tag is somebody else's — a browser's download in progress in +# a shared folder, a copy the operator is making — and is never touched. +_OWN_PART = re.compile(r"\.[0-9a-f]{8}" + re.escape(PART_SUFFIX) + r"$") + + +def is_own_part(path: Path) -> bool: + return bool(_OWN_PART.search(path.name)) + # How long a `.part` with no upload behind it is kept before it is deleted. # # Generous on purpose. The cost of waiting is disk; the cost of being wrong is @@ -107,6 +129,15 @@ class PartialUploads: def drop(self, user_id: str, rel_dir: str, filename: str) -> Partial | None: return self._by_key.pop((user_id, rel_dir, filename), None) + def reserved_names(self, rel_dir: str) -> set[str]: + """The final names uploads in flight into `rel_dir` will take. + + None of them exists on disk yet, so a name check that looked only at + the directory would hand the same name to a second upload. + """ + return {state.stored_name for (_u, d, _f), state in self._by_key.items() + if d == rel_dir} + def __len__(self) -> int: return len(self._by_key) @@ -145,7 +176,7 @@ def orphaned_parts(candidates: Iterable[tuple[Path, float]], """ doomed: list[Path] = [] for path, mtime in candidates: - if path.suffix != PART_SUFFIX: + if not is_own_part(path): continue if path in live: continue 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 |