diff options
| -rw-r--r-- | docs/MESHBAY_DESIGN.md | 14 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/linkpreview.py | 157 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py | 2 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_linkpreview.py | 94 |
4 files changed, 207 insertions, 60 deletions
diff --git a/docs/MESHBAY_DESIGN.md b/docs/MESHBAY_DESIGN.md index 38977af..f1473f0 100644 --- a/docs/MESHBAY_DESIGN.md +++ b/docs/MESHBAY_DESIGN.md @@ -1755,10 +1755,14 @@ card text lives in a bounded in-memory TTL cache; the image rides the same blake3-keyed store as any other thumbnail. **The new surface is SSRF**, because the URL is a member's choice and it triggers an outbound request from the operator's machine: http(s) only, no credentials, a port allowlist, every resolved address -must be globally routable, redirects followed by hand so each hop is re-checked, -and the address the connection landed on re-checked before the body is read (the -request itself has been sent by then, so this refuses the answer, not the -request). The body is read as a stream and stops at its cap — 512 KiB of page, +must be globally routable, and redirects followed by hand so each hop is +re-checked. **The socket is opened to the address that was checked**: the name is +resolved once, off the event loop, every answer is checked, and the connection +goes to that IP literal while TLS verifies the certificate for the name. Checking +one resolution and letting the HTTP client make another would let a name answer +clean and then with a LAN address, and the request would leave before anything +looked; resolving on the event loop would stall every group on a slow name. No +proxy from the environment is used, since a proxy resolves the name itself. The body is read as a stream and stops at its cap — 512 KiB of page, 2 MiB of image — counted after decompression, so a small compressed response is bounded like any other; an image declared larger than its cap is not read, and the whole fetch has a 15-second deadline. Previews are rate-limited per @@ -3456,7 +3460,7 @@ be understood, not so the incident can be retold. | **M2a** | `sender_id` comes from the authenticated session (**NS6**). Now guaranteed by there being **one** chat implementation: QUIC does not carry chat at all (§5.1) | | **M2b** | Chat broadcast is per group (**H1**), on the one transport that carries chat | | **M2c** | No transport runs a synchronous media process on the event loop, and every one is capped | -| **M3** *(third review)* | Link-preview SSRF is gated: rate limit, port allowlist, globally-routable check, per-hop re-check, connect-address re-check, size guard (§6.5) | +| **M3** *(third review)* | Link-preview SSRF is gated: rate limit, port allowlist, globally-routable check, per-hop re-check, connection pinned to the checked address, size guard (§6.5) | | **M4** *(third review)* | Federation binds a pushed row's source to the signer, checks the token audience, caps the push, rejects replays, and scopes revocation to the peer's own entries (§7.6) | | **M5** *(third review)* | A CSP and security headers apply to the hub-served application, verified against the running app — a mis-tuned CSP shows as a blank page | | **M6** *(third review)* | **Withdrawn.** It misread the node registering a hub membership during the CLI invite flow — which is deliberate — as authorization drift | 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/transport/webrtc/chat.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py index 26ec27c..5ef153e 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py @@ -598,7 +598,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/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" |