aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
-rw-r--r--docs/MESHBAY_DESIGN.md14
-rw-r--r--packages/meshbay-node/src/meshbay_node/linkpreview.py157
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py2
-rw-r--r--packages/meshbay-node/tests/test_linkpreview.py94
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"