aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-30 11:22:24 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-30 11:22:24 +0200
commit69554fac7eba6eef7eb8a1c0111c5b92e7f21256 (patch)
tree87ae908d008159c4237b31dade3f385a629c3c7b /packages/meshbay-node/src/meshbay_node
parente2a487c3cd24589b9d9bd298b79d8efaa9914480 (diff)
downloadmeshbay-69554fac7eba6eef7eb8a1c0111c5b92e7f21256.tar.gz
fix: a member can no longer lock a node, crash it with a link, or stop hub cleanup
- node: only a wrong code counts towards the join lock, now per account (5) as well as node-wide (20), and it is consulted only when a code is tried. Every member reconnecting gets the group key through join_request, so a lock checked before recognition let one member refuse it to everyone. - node: link previews read the body as a stream and stop at the cap, counted on decoded bytes; a declared oversized image is not read; 15 s total deadline; image decoding off the loop. `client.get` had buffered the whole (decompressed) response before the caps looked at it. - hub: the daily purge of never-verified accounts detaches their IP-log rows (keeping the name) and clears every other reference first, and each cleanup step runs on its own. On PostgreSQL the bare DELETE violated the ip_logs foreign key and stopped every purge behind it for good. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/linkpreview.py112
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py74
2 files changed, 137 insertions, 49 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/linkpreview.py b/packages/meshbay-node/src/meshbay_node/linkpreview.py
index b223aea..6e1e618 100644
--- a/packages/meshbay-node/src/meshbay_node/linkpreview.py
+++ b/packages/meshbay-node/src/meshbay_node/linkpreview.py
@@ -30,6 +30,7 @@ OG image rides the existing `media_cache` thumb store (same as a poster).
from __future__ import annotations
+import asyncio
import ipaddress
import logging
import socket
@@ -42,6 +43,9 @@ import httpx
log = logging.getLogger(__name__)
_TIMEOUT = 5.0
+# The whole fetch, redirects and body included. `_TIMEOUT` is per network
+# operation, so a server that keeps sending slowly never trips it on its own.
+_TOTAL_DEADLINE = 15.0
_MAX_REDIRECTS = 3
_MAX_HTML_BYTES = 512 * 1024
_MAX_IMAGE_BYTES = 2 * 1024 * 1024
@@ -171,19 +175,56 @@ def _reject_if_rebound(resp: httpx.Response) -> None:
async def _get(client: httpx.AsyncClient, url: str) -> httpx.Response:
- """One GET with manual, re-validated redirects."""
+ """
+ One GET with manual, re-validated redirects, **body not read**.
+
+ The caller reads it through `_read_capped` and must close it. `client.get`
+ is not usable here: it reads and decodes the whole body before returning,
+ so the size caps applied afterwards bounded nothing — a page, an image or a
+ 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)
for _ in range(_MAX_REDIRECTS + 1):
- resp = await client.get(current, headers={"User-Agent": _UA},
- follow_redirects=False)
- _reject_if_rebound(resp)
+ 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:
- current = safe_url(urljoin(current, resp.headers["location"]))
+ location = resp.headers["location"]
+ await resp.aclose()
+ current = safe_url(urljoin(current, location))
continue
return resp
raise UnsafeURL("too many redirects")
+async def _read_capped(resp: httpx.Response, cap: int) -> tuple[bytes, bool]:
+ """
+ At most `cap` bytes of the *decoded* body, and whether there was more.
+
+ Counted after decoding, so a small compressed response that inflates to
+ gigabytes stops at the cap like any other. A declared length over the cap
+ is not read at all.
+ """
+ try:
+ declared = int(resp.headers.get("content-length", ""))
+ except ValueError:
+ declared = -1
+ if resp.headers.get("content-encoding", "identity") == "identity" \
+ and declared > cap:
+ return b"", True
+ body = bytearray()
+ async for chunk in resp.aiter_bytes():
+ body += chunk
+ if len(body) > cap:
+ return bytes(body[:cap]), True
+ return bytes(body), False
+
+
async def fetch_preview(url: str, *, client: httpx.AsyncClient | None = None) -> dict | None:
"""
Return {url, title, description, site_name, image_url} for a URL, or None
@@ -195,17 +236,25 @@ async def fetch_preview(url: str, *, client: httpx.AsyncClient | None = None) ->
if own:
client = httpx.AsyncClient(timeout=_TIMEOUT, max_redirects=0)
try:
- safe_url(url)
- resp = await _get(client, url)
+ return await asyncio.wait_for(_preview(client, url), _TOTAL_DEADLINE)
+ except (httpx.HTTPError, UnsafeURL, TimeoutError) as e:
+ log.debug("link preview for %s: %s", url[:80], e)
+ return None
+ finally:
+ if own:
+ await client.aclose()
+
+
+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()
if resp.status_code != 200 or ctype not in ("text/html", "application/xhtml+xml"):
return None
-
- body = b""
- async for chunk in resp.aiter_bytes():
- body += chunk
- if len(body) >= _MAX_HTML_BYTES:
- break
+ # A page longer than this is cut, not refused: what a card needs is in
+ # the head.
+ body, _truncated = await _read_capped(resp, _MAX_HTML_BYTES)
final_url = str(resp.url)
parser = _HeadParser()
@@ -237,12 +286,8 @@ async def fetch_preview(url: str, *, client: httpx.AsyncClient | None = None) ->
"site_name": (site_name or "")[:120] or None,
"image_url": image,
}
- except (httpx.HTTPError, UnsafeURL) as e:
- log.debug("link preview for %s: %s", url[:80], e)
- return None
finally:
- if own:
- await client.aclose()
+ await resp.aclose()
async def fetch_image(url: str, *, client: httpx.AsyncClient | None = None) -> bytes | None:
@@ -251,23 +296,30 @@ async def fetch_image(url: str, *, client: httpx.AsyncClient | None = None) -> b
if own:
client = httpx.AsyncClient(timeout=_TIMEOUT, max_redirects=0)
try:
- safe_url(url)
- resp = await _get(client, url)
- ctype = resp.headers.get("content-type", "").split(";")[0].strip().lower()
- if resp.status_code != 200 or not ctype.startswith("image/"):
- return None
- raw = b""
- async for chunk in resp.aiter_bytes():
- raw += chunk
- if len(raw) > _MAX_IMAGE_BYTES:
- return None
- return _downscale(raw)
- except (httpx.HTTPError, UnsafeURL) as e:
+ raw = await asyncio.wait_for(_image_bytes(client, url), _TOTAL_DEADLINE)
+ except (httpx.HTTPError, UnsafeURL, TimeoutError) as e:
log.debug("link preview image %s: %s", url[:80], e)
return None
finally:
if own:
await client.aclose()
+ if raw is None:
+ return None
+ # Decoding an image is CPU work on untrusted bytes; not on the event loop.
+ return await asyncio.to_thread(_downscale, raw)
+
+
+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()
+ if resp.status_code != 200 or not ctype.startswith("image/"):
+ return None
+ raw, too_big = await _read_capped(resp, _MAX_IMAGE_BYTES)
+ return None if too_big else raw
+ finally:
+ await resp.aclose()
def _downscale(raw: bytes) -> bytes | None:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py
index e678a13..f94cade 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py
@@ -40,9 +40,18 @@ _INVITE_ID_RE = re.compile(r"[0-9a-f]{32}")
MAX_JOIN_ATTEMPTS = 5
# Per-connection limits alone would not bind an attacker who can open connections
# at will — and the adversary who can mint tokens for any account is the hub. So
-# failed pairings are also counted node-wide over a window.
+# wrong codes are also counted per account and node-wide over a window.
+#
+# Only a *wrong code* counts, and the lock is consulted only when a code is about
+# to be tried. Every member reconnecting obtains the group key through this very
+# message — no per-member bundle is stored — so a lock checked at the top of it,
+# fed by ordinary refusals such as `code_required`, let any one member refuse the
+# key to everybody on the node for as long as they kept failing. A device the
+# node already pinned, joining without a code, is never subject to it.
MAX_JOIN_FAILURES_WINDOW = 20
+MAX_JOIN_FAILURES_PER_ACCOUNT = 5
JOIN_FAILURE_WINDOW = 600 # seconds
+_JOIN_FAILURE_ACCOUNTS_TRACKED = 1000
class AdmissionMixin:
@@ -124,15 +133,11 @@ class AdmissionMixin:
# ── Pairing and join (H3, M3) ────────────────────────────────────────────
- def _join_refuse(self, reason: str, audit_detail: str = "") -> None:
+ def _join_refuse(self, reason: str, audit_detail: str = "", *,
+ wrong_code: bool = False) -> None:
self._join_attempts += 1
- # Node-wide window, shared across connections: reconnecting must not reset
- # the budget.
- now = time.time()
- failures = [t for t in self._ctx.get("join_failures", [])
- if now - t < JOIN_FAILURE_WINDOW]
- failures.append(now)
- self._ctx["join_failures"] = failures
+ if wrong_code:
+ self._record_wrong_code()
self._audit_join("join_refused", audit_detail or reason)
self._send({
"type": MNP.JOIN_RESULT,
@@ -141,6 +146,41 @@ class AdmissionMixin:
"reason": reason,
})
+ def _record_wrong_code(self) -> None:
+ """Count one wrong code, per account and node-wide, across connections:
+ reconnecting must not reset either budget."""
+ now = time.time()
+ failures = [t for t in self._ctx.get("join_failures", [])
+ if now - t < JOIN_FAILURE_WINDOW]
+ failures.append(now)
+ self._ctx["join_failures"] = failures
+
+ per_account: dict = self._ctx.setdefault("join_failures_by_account", {})
+ if len(per_account) > _JOIN_FAILURE_ACCOUNTS_TRACKED:
+ for uid, times in list(per_account.items()):
+ if not times or now - times[-1] >= JOIN_FAILURE_WINDOW:
+ per_account.pop(uid, None)
+ uid = self._user_id or getattr(self, "_pending_sub", "")
+ mine = [t for t in per_account.get(uid, ()) if now - t < JOIN_FAILURE_WINDOW]
+ mine.append(now)
+ per_account[uid] = mine
+
+ def _code_guessing_locked(self, user_id: str) -> bool:
+ """Whether a code may be tried now. Refuses, and says so, when not."""
+ now = time.time()
+ recent = [t for t in self._ctx.get("join_failures", [])
+ if now - t < JOIN_FAILURE_WINDOW]
+ mine = [t for t in (self._ctx.get("join_failures_by_account") or {})
+ .get(user_id, ()) if now - t < JOIN_FAILURE_WINDOW]
+ if len(mine) >= MAX_JOIN_FAILURES_PER_ACCOUNT:
+ self._audit_join("join_throttled", f"{len(mine)} wrong codes by this account")
+ elif len(recent) >= MAX_JOIN_FAILURES_WINDOW:
+ self._audit_join("join_throttled", f"{len(recent)} wrong codes on this node")
+ else:
+ return False
+ self._send({"type": "error", "detail": "Pairing temporarily locked"})
+ return True
+
def _audit_join(self, event: str, detail: str) -> None:
audit = self._ctx.get("audit_store")
if not audit:
@@ -177,14 +217,6 @@ class AdmissionMixin:
self._send({"type": "error", "detail": "Too many attempts"})
return
- now = time.time()
- recent = [t for t in self._ctx.get("join_failures", [])
- if now - t < JOIN_FAILURE_WINDOW]
- if len(recent) >= MAX_JOIN_FAILURES_WINDOW:
- self._audit_join("join_throttled", f"{len(recent)} failures in window")
- self._send({"type": "error", "detail": "Pairing temporarily locked"})
- return
-
user_id = self._user_id or getattr(self, "_pending_sub", "")
username = self._username or getattr(self, "_pending_username", "")
if not user_id:
@@ -295,9 +327,11 @@ class AdmissionMixin:
if not code:
self._join_refuse("code_required")
return
+ if self._code_guessing_locked(user_id):
+ return
invite = await roster.consume_invite(code, user_id, session_group)
if not invite:
- self._join_refuse("code_invalid")
+ self._join_refuse("code_invalid", wrong_code=True)
return
await roster.set_member(
group_id=invite["group_id"], user_id=user_id,
@@ -337,9 +371,11 @@ class AdmissionMixin:
self._join_refuse("code_required")
return
+ if self._code_guessing_locked(user_id):
+ return
invite = await roster.consume_invite(code, user_id, session_group)
if not invite:
- self._join_refuse("code_invalid")
+ self._join_refuse("code_invalid", wrong_code=True)
return
await self._pin_and_admit(