aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-01 21:51:25 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-01 21:51:25 +0200
commit799d87999c8324564dce5159191532e008dd93d2 (patch)
tree5ff1816f18625dfece9eb67fa06c7b25fdece4f8 /packages
parent8a6294b0412a86f378c6e2e937c28de64a903c91 (diff)
parent1e6db7d23c70b7bd7e1422f09911b3645f0fb2e2 (diff)
downloadmeshbay-799d87999c8324564dce5159191532e008dd93d2.tar.gz
Merge branch 'fix/third-review-h1-h2-m1-m6'
Third security review (docs/third-review.md) plus its remediation. Fixed and verified: - H1 moderator could grant admin / hard-revoke → handler split by field - H2 unauthenticated 2-report global blocklist → auth + distinct reporters + rate limit + refused when public groups are off - M1 registration reCAPTCHA was inert → gate unconditional; the desktop client's CSP allows the widget - M2 QUIC chat/stream handlers lagged WebRTC → brought to parity; the QUIC listener is now off by default ([node] quic_enabled) - M3 link-preview SSRF gaps → rate limit + port allowlist + connect-address re-check + decompression-bomb guard - M4 federated peer over-trust → source bound to the signer, push capped, revocation prunes the peer's own entries, replay rejected - M5 no CSP / security headers on the SPA → middleware; verified against the live app with no violations Withdrawn: - M6 add_group_member accepting node tokens is deliberate (commit 0443cf8, the CLI invite flow). The "fix" broke that flow on the deployed hub and was reverted. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011pG75yGK3NthNfyjH74omG
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-client/src/main.js19
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/admin.py21
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/deps.py17
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/federation.py124
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/groups.py7
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/moderation.py86
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/users.py9
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/webapp.py31
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/app.py17
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/auth-page.js9
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/keyderive.js6
-rw-r--r--packages/meshbay-hub/tests/test_admin.py39
-rw-r--r--packages/meshbay-hub/tests/test_desktop_shell.py42
-rw-r--r--packages/meshbay-hub/tests/test_federation.py179
-rw-r--r--packages/meshbay-hub/tests/test_moderation.py132
-rw-r--r--packages/meshbay-hub/tests/test_node_auth.py23
-rw-r--r--packages/meshbay-hub/tests/test_register_captcha.py58
-rw-r--r--packages/meshbay-hub/tests/test_security_headers.py65
-rw-r--r--packages/meshbay-node/src/meshbay_node/config.py15
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py9
-rw-r--r--packages/meshbay-node/src/meshbay_node/linkpreview.py51
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/quic_server.py144
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py40
-rw-r--r--packages/meshbay-node/tests/test_link_preview_request.py50
-rw-r--r--packages/meshbay-node/tests/test_linkpreview.py40
-rw-r--r--packages/meshbay-node/tests/test_quic_enabled.py53
26 files changed, 1105 insertions, 181 deletions
diff --git a/packages/meshbay-client/src/main.js b/packages/meshbay-client/src/main.js
index 3926c01..e071074 100644
--- a/packages/meshbay-client/src/main.js
+++ b/packages/meshbay-client/src/main.js
@@ -58,15 +58,29 @@ const SCHEME = 'app';
// The hub is reachable under connect-src, for its API and its signaling socket.
// It is deliberately absent from script-src: nothing it returns is executed,
// which is the whole reason this application exists (T3).
+//
+// The one exception is reCAPTCHA, used to gate sign-up (and password reset) the
+// same way it gates them in the browser. Its script comes from www.google.com,
+// its challenge is a www.google.com iframe, and its assets sit on
+// www.gstatic.com. These two hosts — and only these two — are allowed under
+// `script-src`, `frame-src` and `img-src` for that purpose. It is a real, if
+// small, dent in "no third-party code runs here": Google's reCAPTCHA script
+// executes in the renderer. It is accepted deliberately so a native sign-up is
+// gated like a web one without asking the user to do anything extra, and it is
+// the *same* dependency the hub-served SPA already carries. If sign-up ever
+// moves to a proof-of-work challenge, delete RECAPTCHA_SRC and the three
+// directives that spread it, and the widget in auth-page.js with them.
+const RECAPTCHA_SRC = 'https://www.google.com https://www.gstatic.com';
const CSP = [
"default-src 'none'",
- "script-src 'self' 'wasm-unsafe-eval'",
+ `script-src 'self' 'wasm-unsafe-eval' ${RECAPTCHA_SRC}`,
"style-src 'self' 'unsafe-inline'",
- "img-src 'self' data: blob:",
+ `img-src 'self' data: blob: ${RECAPTCHA_SRC}`,
"media-src 'self' blob:",
"font-src 'self'",
"connect-src 'self' https: wss:",
"worker-src 'self'",
+ `frame-src ${RECAPTCHA_SRC}`,
"frame-ancestors 'none'",
"base-uri 'none'",
"form-action 'none'",
@@ -857,6 +871,7 @@ function registerBridge() {
`username = "${username}"`,
'',
'[node]',
+ 'quic_enabled = false # QUIC direct path; no client uses it yet',
'quic_port = 19010',
'ui_port = 18000',
'',
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/admin.py b/packages/meshbay-hub/src/meshbay_hub/api/admin.py
index ab6fa07..4960674 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/admin.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/admin.py
@@ -14,7 +14,7 @@ from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from meshbay_hub.auth import decrypt_email
-from meshbay_hub.api.deps import require_admin, require_moderator
+from meshbay_hub.api.deps import require_admin, require_moderator, user_is_admin
from meshbay_hub.api.revocation import get_connected_node_count, is_node_connected
from meshbay_hub.db.engine import get_db
from meshbay_hub.db.models import Group, GroupMember, IPLog, Node, User
@@ -196,6 +196,25 @@ async def admin_patch_user(
if user.id == current_user.id:
raise HTTPException(status_code=400, detail="Cannot modify your own account")
+ # A moderator suspends and restores accounts — reversible content moderation.
+ # Changing what someone *is* (their role), and the one irreversible status
+ # (`revoked`, which is signed and broadcast to every node), are administrative.
+ # Without this split a moderator could promote an accomplice to admin, or
+ # revoke every admin, entirely from the moderation role. `admin_delete_user`
+ # already draws this exact line for the same reason.
+ privileged = body.role is not None or body.status == "revoked"
+ if privileged and not user_is_admin(current_user):
+ raise HTTPException(
+ status_code=403,
+ detail="Changing a role, or revoking an account, requires admin rights")
+
+ # An admin's account is not a moderator's to touch at all — not their role,
+ # not their status.
+ if user_is_admin(user) and not user_is_admin(current_user):
+ raise HTTPException(
+ status_code=403,
+ detail="Only an admin can change another admin's account")
+
from meshbay_hub.api.notifications import create_notification
if body.role is not None:
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/deps.py b/packages/meshbay-hub/src/meshbay_hub/api/deps.py
index 1bf57a4..907a481 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/deps.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/deps.py
@@ -72,11 +72,21 @@ async def require_user_scope(
return current_user
+def user_is_admin(user: User) -> bool:
+ """Admin by DB role or by the config allow-list. Use inside a handler that
+ already depends on `require_moderator` but has to draw the admin line for
+ one field (see `admin_patch_user`)."""
+ return user.role == "admin" or user.username in _admin_usernames
+
+
+def user_is_moderator(user: User) -> bool:
+ return user.role in ("moderator", "admin") or user.username in _admin_usernames
+
+
async def require_moderator(
current_user: User = Depends(get_current_user),
) -> User:
- if current_user.role not in ("moderator", "admin") \
- and current_user.username not in _admin_usernames:
+ if not user_is_moderator(current_user):
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN,
detail="Moderator access required")
return current_user
@@ -85,8 +95,7 @@ async def require_moderator(
async def require_admin(
current_user: User = Depends(get_current_user),
) -> User:
- if current_user.role != "admin" \
- and current_user.username not in _admin_usernames:
+ if not user_is_admin(current_user):
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN,
detail="Admin access required")
return current_user
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/federation.py b/packages/meshbay-hub/src/meshbay_hub/api/federation.py
index 7f1262d..b1e0e30 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/federation.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/federation.py
@@ -27,7 +27,7 @@ import uuid
import jwt
from fastapi import APIRouter, Depends, HTTPException, Header
from pydantic import BaseModel
-from sqlalchemy import select
+from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from meshbay_common import MHP_VERSION
@@ -41,6 +41,17 @@ log = logging.getLogger(__name__)
router = APIRouter(prefix="/mhp", tags=["federation"])
+# One push may not dump the world, and one peer may not fill the table.
+MAX_FEDERATED_GROUPS_PER_PUSH = 500
+MAX_FEDERATED_GROUPS_PER_PEER = 2000
+
+# Seen `jti` values for the state-changing MHP endpoints, pruned lazily. The
+# sending side that would set an `aud` claim is unbuilt, so audience binding is
+# not available; this stops a captured POST /mhp/directory or /mhp/revoke from
+# being replayed inside the token's short TTL. GET /mhp/directory is idempotent
+# and not covered.
+_seen_mhp_jti: dict[str, float] = {}
+
def _issue_mhp_token(target_hub_id: str) -> str:
"""Issue a short-lived JWT for authenticating to a peer hub."""
@@ -57,9 +68,15 @@ def _issue_mhp_token(target_hub_id: str) -> str:
async def _verify_mhp_token(
- token: str, db: AsyncSession, expected_aud: str | None = None,
+ token: str, db: AsyncSession, *, single_use: bool = False,
) -> dict:
- """Verify a JWT from a peer hub using DB-stored public key."""
+ """
+ Verify a JWT from a peer hub against its DB-stored public key and return the
+ payload.
+
+ `single_use=True` (the state-changing endpoints) additionally rejects a
+ replayed `jti` within the token's lifetime.
+ """
unverified = jwt.decode(token, options={"verify_signature": False})
sender_id = unverified.get("iss")
@@ -67,15 +84,22 @@ async def _verify_mhp_token(
if not peer:
raise PermissionError(f"Unknown hub: {sender_id!r}. Register as peer first.")
- options = {}
- if expected_aud:
- options["audience"] = expected_aud
-
decoded = jwt.decode(
token, peer.pk_hub_pem.encode(),
algorithms=["EdDSA"],
- options=options,
+ options={"require": ["exp", "iss"]},
)
+
+ if single_use:
+ now = time.time()
+ for j, exp in list(_seen_mhp_jti.items()):
+ if exp < now:
+ _seen_mhp_jti.pop(j, None)
+ jti = decoded.get("jti", "")
+ if not jti or jti in _seen_mhp_jti:
+ raise PermissionError("MHP token replay")
+ _seen_mhp_jti[jti] = float(decoded.get("exp", now + 300))
+
return decoded
@@ -142,30 +166,50 @@ async def receive_directory(
db: AsyncSession = Depends(get_db),
):
try:
- await _verify_mhp_token(authorization.removeprefix("Bearer "), db)
+ payload = await _verify_mhp_token(
+ authorization.removeprefix("Bearer "), db, single_use=True)
except Exception as e:
raise HTTPException(status_code=401, detail=str(e))
+ # `source_hub` is the signer of this request, never `body.hub_id` — a peer
+ # does not get to relay or spoof a third hub's groups into our directory.
+ sender = payload["iss"]
+ if len(body.groups) > MAX_FEDERATED_GROUPS_PER_PUSH:
+ raise HTTPException(status_code=413, detail="Too many groups in one push")
+
from datetime import datetime, timezone
now = datetime.now(timezone.utc)
+ have = await db.scalar(
+ select(func.count()).select_from(FederatedGroup)
+ .where(FederatedGroup.source_hub == sender)) or 0
+
count = 0
for g in body.groups:
- existing = await db.get(FederatedGroup, g["id"])
- if existing:
- existing.name = g.get("name", existing.name)
- existing.join_policy = g.get("join_policy", existing.join_policy)
- existing.updated_at = now
+ gid = str(g.get("id", ""))[:36]
+ name = str(g.get("name", ""))[:128]
+ jp = g.get("join_policy", "invite")
+ if not gid or jp not in ("invite", "open"):
+ continue
+ # A federated id must never shadow a real local group.
+ if await db.get(Group, gid):
+ log.warning("Federated id %s collides with a local group — skipped", gid[:8])
+ continue
+ row = await db.get(FederatedGroup, gid)
+ if row:
+ if row.source_hub != sender:
+ continue # only the hub that advertised it may update it
+ row.name = name or row.name
+ row.join_policy = jp
+ row.updated_at = now
else:
+ if have + count >= MAX_FEDERATED_GROUPS_PER_PEER:
+ break
db.add(FederatedGroup(
- id=g["id"],
- name=g.get("name", ""),
- source_hub=body.hub_id,
- join_policy=g.get("join_policy", "invite"),
- ))
+ id=gid, name=name, source_hub=sender, join_policy=jp))
count += 1
await db.commit()
- log.info("Persisted %d groups from hub %s", count, body.hub_id[:16])
- return {"accepted": count, "from_hub": body.hub_id}
+ log.info("Persisted %d groups from hub %s", count, sender[:16])
+ return {"accepted": count, "from_hub": sender}
# ── Revocation propagation ────────────────────────────────────────────────────
@@ -179,15 +223,43 @@ async def receive_revocation(
authorization: str = Header(...),
db: AsyncSession = Depends(get_db),
):
+ """
+ Act on a revocation from a peer hub.
+
+ This does **not** reach local nodes: nothing here hosts a federated group,
+ and a local node would reject a token signed by another hub's key anyway
+ (that path was a silent no-op). What a peer may legitimately revoke is a
+ group *it advertised to us* — so this prunes our copy of the peer's
+ directory. Users are per-hub; a peer does not get to revoke ours.
+ """
try:
- await _verify_mhp_token(authorization.removeprefix("Bearer "), db)
+ payload = await _verify_mhp_token(
+ authorization.removeprefix("Bearer "), db, single_use=True)
except Exception as e:
raise HTTPException(status_code=401, detail=str(e))
- from meshbay_hub.api.revocation import broadcast_revocation
- sent = await broadcast_revocation(body.token)
- log.info("Propagated revocation to %d local nodes", sent)
- return {"propagated_to": sent}
+ sender = payload["iss"]
+ peer = await db.get(HubPeer, sender)
+ try:
+ inner = jwt.decode(
+ body.token, peer.pk_hub_pem.encode(), algorithms=["EdDSA"],
+ options={"verify_exp": False})
+ except Exception as e:
+ raise HTTPException(status_code=400, detail=f"Bad revocation token: {e}")
+
+ if inner.get("type") != "revocation" or inner.get("target") != "group":
+ return {"pruned": 0, "note": "federation may only revoke groups it advertised"}
+
+ target_id = inner.get("target_id", "")
+ row = await db.get(FederatedGroup, target_id)
+ pruned = 0
+ if row and row.source_hub == sender:
+ await db.delete(row)
+ await db.commit()
+ pruned = 1
+ log.info("Federated group %s revoked by %s (pruned=%d)",
+ target_id[:8], sender[:16], pruned)
+ return {"pruned": pruned}
# ── Peer management (admin) ───────────────────────────────────────────────────
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
index fdb0444..07d3ec0 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/groups.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
@@ -589,6 +589,13 @@ async def add_group_member(
current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
+ # `get_current_user`, not `require_user_scope`: the node calls this after a
+ # CLI `member invite` so the group becomes visible in the invitee's SPA
+ # (commit 0443cf8). The node authenticates with a node-scoped token, and the
+ # `group.admin_id == current_user.id` check below is the real guard — a node
+ # can only touch its own operator's groups, adding an already-registered
+ # account. (Third-review M6 proposed tightening this to `require_user_scope`;
+ # that broke the CLI invite flow and was reverted — see the review doc.)
group = await db.get(Group, group_id)
if not group:
raise HTTPException(status_code=404, detail="Group not found")
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py
index 853f255..0939eec 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py
@@ -1,13 +1,16 @@
"""
MeshBay Hub — moderation endpoints.
-Public reporting flow:
- POST /v1/reports — report a content hash (no auth required)
+Reporting flow:
+ POST /v1/reports — report a content hash (sign-in required)
- Thresholds:
- 1st report → logged, node admin notified (future: push notification)
- 2nd report → content hash added to blocklist automatically
- 3rd+ report → logged as repeat offense (escalation for human review)
+ Thresholds (counted as DISTINCT reporting accounts, not raw rows):
+ < AUTO_BLOCK_THRESHOLD distinct reporters → logged
+ >= AUTO_BLOCK_THRESHOLD distinct reporters → hash added to the blocklist
+
+ The flow only runs while the hub brokers public content: with public groups
+ switched off instance-wide there is nothing here to serve a reported hash from,
+ so it is refused rather than left open as an unauthenticated write surface.
Admin endpoints:
GET /v1/admin/blocklist — list blocked hashes
@@ -27,7 +30,9 @@ from pydantic import BaseModel
from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
+from meshbay_hub import hub_settings
from meshbay_hub.api.deps import get_current_user, require_admin
+from meshbay_hub.api.middleware import limiter
from meshbay_hub.api.netutil import client_ip
from meshbay_hub.db.engine import get_db
from meshbay_hub.db.models import ContentBlocklist, ContentReport, User
@@ -36,7 +41,11 @@ log = logging.getLogger(__name__)
router = APIRouter(tags=["moderation"])
-AUTO_BLOCK_THRESHOLD = 2 # reports before automatic block
+# Distinct reporting accounts before a hash is auto-blocked. Kept low for a
+# responsive community signal, but note it is only as strong as account
+# creation: while a bot can register freely (see the reCAPTCHA gap), the real
+# control is the admin reviewing `GET /v1/admin/blocklist` and the audit log.
+AUTO_BLOCK_THRESHOLD = 3
# ── Models ────────────────────────────────────────────────────────────────────
@@ -56,34 +65,58 @@ class BlocklistAddRequest(BaseModel):
# ── Public endpoints ──────────────────────────────────────────────────────────
@router.post("/v1/reports", status_code=201)
+@limiter.limit("10/hour")
async def report_content(
body: ReportRequest,
request: Request,
+ current_user: User = Depends(get_current_user),
db: AsyncSession = Depends(get_db),
):
- """Report a content hash. No authentication required."""
+ """
+ Report a public content hash for moderation.
+
+ Sign-in is required. It used to be anonymous, which made it a censorship
+ primitive: two unauthenticated POSTs naming any blake3 id auto-added it to
+ the blocklist that nodes enforce, network-wide, with manual admin removal the
+ only undo. The threshold now counts *distinct reporting accounts*, one vote
+ per account per hash.
+
+ Refused entirely when the hub has public groups switched off: nothing here
+ brokers public content then, nothing syncs the blocklist, and an open write
+ endpoint would only be abuse surface.
+ """
+ if not await hub_settings.public_groups_allowed(db):
+ raise HTTPException(
+ status_code=403,
+ detail="This hub does not broker public content, so there is nothing to report here.")
+
if len(body.content_hash) != 64 or not all(c in "0123456789abcdef" for c in body.content_hash):
raise HTTPException(status_code=422, detail="content_hash must be 64 hex chars (blake3)")
- ip = client_ip(request)
+ # One vote per account per hash — a single reporter must not be able to walk
+ # the threshold up on their own by posting repeatedly.
+ already = await db.scalar(
+ select(ContentReport.id).where(
+ ContentReport.content_hash == body.content_hash,
+ ContentReport.reporter_id == current_user.id))
- # Count existing reports for this hash
- count_result = await db.execute(
- select(func.count()).where(ContentReport.content_hash == body.content_hash))
- count = count_result.scalar_one()
+ if not already:
+ db.add(ContentReport(
+ content_hash=body.content_hash,
+ reporter_id=current_user.id,
+ group_id=body.group_id,
+ reason=body.reason,
+ detail=body.detail,
+ ip_address=client_ip(request),
+ ))
+ await db.flush()
- report = ContentReport(
- content_hash=body.content_hash,
- group_id=body.group_id,
- reason=body.reason,
- detail=body.detail,
- ip_address=ip,
- )
- db.add(report)
+ distinct_reporters = await db.scalar(
+ select(func.count(func.distinct(ContentReport.reporter_id)))
+ .where(ContentReport.content_hash == body.content_hash)) or 0
- action = "logged"
- if count + 1 >= AUTO_BLOCK_THRESHOLD:
- # Check if already blocked
+ action = "already_reported" if already else "logged"
+ if distinct_reporters >= AUTO_BLOCK_THRESHOLD:
existing = await db.get(ContentBlocklist, body.content_hash)
if not existing:
db.add(ContentBlocklist(
@@ -92,13 +125,14 @@ async def report_content(
added_by="auto",
))
action = "auto_blocked"
- log.warning("Content auto-blocked after %d reports: %s", count + 1, body.content_hash[:16])
+ log.warning("Content auto-blocked after %d distinct reporters: %s",
+ distinct_reporters, body.content_hash[:16])
await db.commit()
return {
"status": action,
"content_hash": body.content_hash,
- "report_count": count + 1,
+ "report_count": distinct_reporters,
"threshold": AUTO_BLOCK_THRESHOLD,
}
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py
index 559cfa6..9b70b59 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/users.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py
@@ -150,8 +150,13 @@ async def register(
return {"user_id": found.id, "email_verification_required": True}
raise HTTPException(status_code=409, detail="Username already taken")
- # Captcha gate — web path only (native clients send auth_key)
- if _cfg and _cfg.captcha.enabled and not body.auth_key:
+ # Captcha gate — every fresh registration when a captcha is configured, with
+ # no client carve-out. The earlier `and not body.auth_key` exempted anything
+ # that sent an `auth_key`, which is *every* real client (the browser sends it
+ # too, from the password split) — so the check was off for everyone, and a
+ # bot skipped it by sending the field. The desktop client is Chromium and
+ # renders the same widget, so it has no need of an exemption either.
+ if _cfg and _cfg.captcha.enabled:
await _verify_captcha_or_raise(body.captcha_token, request)
# Email uniqueness (only active or pending accounts)
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/webapp.py b/packages/meshbay-hub/src/meshbay_hub/api/webapp.py
index 0821809..6dfd3ed 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/webapp.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/webapp.py
@@ -22,7 +22,8 @@ STATIC_DIR = Path(__file__).parent.parent / "static"
router = APIRouter(tags=["webapp"])
# Assets the shell pulls in, in load order. Everything else is imported by
-# app.js and rides on the same query string via window.__MB_ASSET_V.
+# app.js from a relative path, which inherits the `/a/<hash>/` prefix the shell
+# loaded app.js under — so the whole module graph moves together.
# Every module the page loads. A file missing from here is a file whose change
# does not move the URL, so a browser holding the old one never asks for it —
# which is the failure this list exists to prevent, and it is silent.
@@ -77,6 +78,33 @@ ASSET_V = _asset_version()
_NO_STORE = {"Cache-Control": "no-store"}
+# Content-Security-Policy for the whole hub, applied by a middleware in app.py.
+#
+# This is the *same* policy the desktop client's protocol handler already sends
+# for these exact UI files (`meshbay-client/src/main.js`), plus the two reCAPTCHA
+# hosts the sign-up widget loads its script, challenge iframe and images from.
+# `'unsafe-inline'` is style-only — htm/preact set inline `style=` attributes
+# everywhere; nothing inline executes, and the shell below carries no inline
+# `<script>`. `'wasm-unsafe-eval'` is required for the Argon2id WASM. The hub's
+# own origin is deliberately absent from `script-src`: a response it returns is
+# never executed, which is the point of T3.
+_RECAPTCHA_SRC = "https://www.google.com https://www.gstatic.com"
+CSP = "; ".join([
+ "default-src 'none'",
+ f"script-src 'self' 'wasm-unsafe-eval' {_RECAPTCHA_SRC}",
+ "style-src 'self' 'unsafe-inline'",
+ f"img-src 'self' data: blob: {_RECAPTCHA_SRC}",
+ "media-src 'self' blob:",
+ "font-src 'self'",
+ "connect-src 'self' https: wss:",
+ "worker-src 'self'",
+ f"frame-src {_RECAPTCHA_SRC}",
+ "frame-ancestors 'none'",
+ "base-uri 'none'",
+ "form-action 'none'",
+])
+
+
@router.get("/app", response_class=HTMLResponse)
async def app_root():
return HTMLResponse(_HTML, headers=_NO_STORE)
@@ -109,7 +137,6 @@ _HTML = """\
and app.js's own relative imports inherit the prefix, which is the only
way the module graph is guaranteed not to be a mixture of two builds.
See _asset_version() and VersionedStatics. -->
- <script>window.__MB_ASSET_V = "{v}";</script>
<!-- Argon2id (WebAssembly, inlined) — WebCrypto has no memory-hard KDF, and the
keypair bundle needs one: it is protected by the passphrase alone and sits
on every node its owner joins (C4). Vendored, see static/vendor/PROVENANCE.md -->
diff --git a/packages/meshbay-hub/src/meshbay_hub/app.py b/packages/meshbay-hub/src/meshbay_hub/app.py
index 76d7ec0..2daa55b 100644
--- a/packages/meshbay-hub/src/meshbay_hub/app.py
+++ b/packages/meshbay-hub/src/meshbay_hub/app.py
@@ -35,7 +35,7 @@ from meshbay_hub.api.relay import router as relay_router
from meshbay_hub.api.signaling import router as signaling_router
from meshbay_hub.api.admin import router as admin_router
from meshbay_hub.api.notifications import router as notifications_router
-from meshbay_hub.api.webapp import router as webapp_router, STATIC_DIR, ASSET_V
+from meshbay_hub.api.webapp import router as webapp_router, STATIC_DIR, ASSET_V, CSP
from meshbay_hub.api.middleware import limiter
@@ -142,6 +142,21 @@ def create_app(cfg: HubConfig | None = None) -> FastAPI:
app.state.limiter = limiter
app.add_exception_handler(RateLimitExceeded, _rate_limit_exceeded_handler)
+ @app.middleware("http")
+ async def _security_headers(request, call_next):
+ """
+ Second-review L5, third-review M5: the SPA shell and its assets went out
+ with no CSP and no other protective headers. This adds them everywhere —
+ `webapp.CSP` is the same policy the desktop client already enforces on
+ these exact files. `setdefault` so a route that sets its own wins.
+ """
+ response = await call_next(request)
+ response.headers.setdefault("Content-Security-Policy", CSP)
+ response.headers.setdefault("X-Content-Type-Options", "nosniff")
+ response.headers.setdefault("Referrer-Policy", "strict-origin-when-cross-origin")
+ response.headers.setdefault("X-Frame-Options", "DENY")
+ return response
+
# Routers (webapp last — catches / before API routes)
app.include_router(hub_router)
app.include_router(users_router)
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/auth-page.js b/packages/meshbay-hub/src/meshbay_hub/static/auth-page.js
index df08bc4..4c00137 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/auth-page.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/auth-page.js
@@ -237,8 +237,11 @@ export function RegisterPage() {
const rk = window.MeshBayKeys.generateRecoveryKey();
// `name` (trimmed), not the raw field: the hub stores the trimmed
// username and every key derivation must fold in the same string.
+ // `captcha.token` rides along — the submit button is already disabled
+ // until it is set when a captcha is configured (see the form below).
await window.MeshBayKeys.registerUser(
- name, email, password, emailRecovery ? rk.mnemonic : null);
+ name, email, password, emailRecovery ? rk.mnemonic : null,
+ captcha.token);
setRecoveryMnemonic(rk.mnemonic);
session.recoveryKey =
await window.MeshBayKeys.deriveRecoveryKey(rk.mnemonic, name);
@@ -257,6 +260,10 @@ export function RegisterPage() {
}
} catch (err) {
setError(err.message);
+ // A reCAPTCHA token is single-use: after a failed attempt (name taken,
+ // e-mail in use…) it is spent, so clear it and make the user solve a
+ // fresh one before the next try. No-op when no captcha is configured.
+ captcha.reset();
} finally {
setLoading(false);
}
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js
index 0aaa6a5..a540a94 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js
@@ -276,7 +276,7 @@ async function decryptBundle(bundleB64, password, username) {
*
* Returns the raw private keys for immediate use after registration.
*/
-async function registerUser(username, email, password, recoveryMnemonic) {
+async function registerUser(username, email, password, recoveryMnemonic, captchaToken) {
// No keypair here any more. Identity keys are per node: one is generated the
// first time this account joins a given node, encrypted under the passphrase,
// and left with that node. So an operator who cracks what sits on their own
@@ -291,6 +291,10 @@ async function registerUser(username, email, password, recoveryMnemonic) {
// appends it to the verification e-mail and stores it nowhere
// (docs/auth-confirm.md §4.4). Omitted when they chose to save it themselves.
if (recoveryMnemonic) payload.recovery_key = recoveryMnemonic;
+ // reCAPTCHA response, when the hub has a captcha configured. The widget lives
+ // in RegisterPage (auth-page.js); this function just forwards its token. A
+ // hub with no captcha configured sends nothing and the server does not check.
+ if (captchaToken) payload.captcha_token = captchaToken;
const resp = await hubCall('/v1/users/register', {
method: 'POST',
diff --git a/packages/meshbay-hub/tests/test_admin.py b/packages/meshbay-hub/tests/test_admin.py
index 51b233a..ad48487 100644
--- a/packages/meshbay-hub/tests/test_admin.py
+++ b/packages/meshbay-hub/tests/test_admin.py
@@ -5,7 +5,6 @@ Integration tests for the admin/moderation panel API.
import pytest
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey
-
from meshbay_common.crypto import pk_to_b64
from meshbay_hub.api.deps import set_admin_usernames
@@ -162,6 +161,44 @@ async def test_admin_change_role(client):
@pytest.mark.asyncio
+async def test_moderator_cannot_change_roles_or_revoke(client):
+ """A moderator suspends and restores (reversible); it cannot promote anyone
+ or hard-revoke, which would be a path from the moderation role to full
+ instance control."""
+ _, admin_token = await _setup_admin(client, "boss")
+ mod_id = await _register(client, "moduser")
+ await client.patch(f"/v1/admin/users/{mod_id}", json={"role": "moderator"},
+ headers={"Authorization": f"Bearer {admin_token}"})
+ mod_token = await _login(client, "moduser")
+ mod_h = {"Authorization": f"Bearer {mod_token}"}
+
+ victim = await _register(client, "victim", email="v@x.com")
+
+ # No promoting an accomplice.
+ r = await client.patch(f"/v1/admin/users/{victim}", json={"role": "admin"},
+ headers=mod_h)
+ assert r.status_code == 403
+
+ # No hard revocation.
+ r = await client.patch(f"/v1/admin/users/{victim}", json={"status": "revoked"},
+ headers=mod_h)
+ assert r.status_code == 403
+
+ # No touching an admin's account.
+ admin2 = await _register(client, "admin2", email="a2@x.com")
+ await client.patch(f"/v1/admin/users/{admin2}", json={"role": "admin"},
+ headers={"Authorization": f"Bearer {admin_token}"})
+ r = await client.patch(f"/v1/admin/users/{admin2}", json={"status": "suspended"},
+ headers=mod_h)
+ assert r.status_code == 403
+
+ # Suspending a plain user is still fine.
+ r = await client.patch(f"/v1/admin/users/{victim}", json={"status": "suspended"},
+ headers=mod_h)
+ assert r.status_code == 200
+
+
+@pytest.mark.asyncio
async def test_admin_cannot_modify_self(client):
admin_id, token = await _setup_admin(client)
r = await client.patch(f"/v1/admin/users/{admin_id}",
diff --git a/packages/meshbay-hub/tests/test_desktop_shell.py b/packages/meshbay-hub/tests/test_desktop_shell.py
index 36b261e..b82804e 100644
--- a/packages/meshbay-hub/tests/test_desktop_shell.py
+++ b/packages/meshbay-hub/tests/test_desktop_shell.py
@@ -175,11 +175,17 @@ def _policy() -> str:
"""
import re
source = _main()
+ # The array mixes plain strings and one `${RECAPTCHA_SRC}` template literal;
+ # resolve the constant so every directive reads as plain text.
+ rec = re.search(r"const RECAPTCHA_SRC = '([^']*)'", source)
match = re.search(r"const CSP = \[(.*?)\]\.join", source, re.S)
assert match, "no CSP constant in the main process"
+ body = match.group(1)
+ if rec:
+ body = body.replace("${RECAPTCHA_SRC}", rec.group(1))
return "; ".join(
- line.strip().strip('",').strip('"')
- for line in match.group(1).splitlines() if line.strip())
+ line.strip().strip('`",').strip('`"')
+ for line in body.splitlines() if line.strip())
def _directive(name: str) -> str:
@@ -193,18 +199,46 @@ def _directive(name: str) -> str:
def test_the_hub_is_reachable_but_never_executable():
"""
connect-src allows the hub's API and its signaling socket. script-src does
- not include it: nothing the hub returns is ever executed.
+ not: nothing the hub returns is ever executed. The only script sources are
+ 'self', the wasm eval token, and the two reCAPTCHA hosts (see the next
+ test) — never a bare `https:` scheme, which would let the hub's own origin
+ serve script.
"""
connect = _directive("connect-src")
assert "https:" in connect and "wss:" in connect
script = _directive("script-src")
assert script, "no script-src directive"
- assert "https:" not in script, "the hub can serve script under this policy"
+ sources = script.split()[1:] # drop the "script-src" keyword itself
+ allowed = {
+ "'self'", "'wasm-unsafe-eval'",
+ "https://www.google.com", "https://www.gstatic.com",
+ }
+ assert set(sources) <= allowed, \
+ f"unexpected script-src source: {set(sources) - allowed}"
+ assert "https:" not in sources, "a bare https: scheme lets the hub serve script"
assert "'unsafe-eval'" not in script.replace("'wasm-unsafe-eval'", "")
assert "default-src 'none'" in _policy()
+def test_recaptcha_is_the_only_third_party_and_stays_scoped_to_it():
+ """
+ reCAPTCHA gates sign-up in the app the same way it does in the browser.
+ www.google.com and www.gstatic.com are allowed under script-src, frame-src
+ and img-src for that — and no other external origin appears anywhere in the
+ policy. Remove this expectation only alongside the reCAPTCHA widget.
+ """
+ hosts = {"https://www.google.com", "https://www.gstatic.com"}
+ for directive in ("script-src", "frame-src", "img-src"):
+ srcs = set(_directive(directive).split()[1:])
+ assert hosts <= srcs, f"{directive} is missing a reCAPTCHA host"
+
+ for part in _policy().split(";"):
+ for tok in part.strip().split()[1:]:
+ if tok.startswith(("http://", "https://")):
+ assert tok in hosts, f"unexpected external origin in CSP: {tok}"
+
+
# ── The bridge ──────────────────────────────────────────────────────────────
def test_the_bridge_is_the_only_way_in():
diff --git a/packages/meshbay-hub/tests/test_federation.py b/packages/meshbay-hub/tests/test_federation.py
new file mode 100644
index 0000000..6035b0c
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_federation.py
@@ -0,0 +1,179 @@
+"""
+MHP federation — what a registered peer hub may and may not do.
+
+A peer is trusted enough to advertise its own public groups into our directory
+and to withdraw them. It is not trusted to speak for a third hub, to shadow a
+local group, to revoke our users, or to replay a state-changing request.
+"""
+
+import base64
+import hashlib
+import time
+import uuid
+
+import jwt
+import pytest
+from cryptography.hazmat.primitives import serialization
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+from meshbay_hub.api.deps import set_admin_usernames
+
+
+def _auth_key(password: str, username: str) -> str:
+ salt = hashlib.sha256(f"meshbay:auth:v1:{username}".encode()).digest()
+ return base64.b64encode(
+ hashlib.pbkdf2_hmac("sha512", password.encode(), salt, 600_000, 32)).decode()
+
+
+async def _admin(client, username="root"):
+ pw = "a-long-enough-passphrase"
+ await client.post("/v1/users/register", json={
+ "username": username, "email": f"{username}@example.com",
+ "auth_key": _auth_key(pw, username)})
+ set_admin_usernames([username])
+ r = await client.post("/v1/users/login", json={
+ "username": username, "auth_key": _auth_key(pw, username)})
+ return {"Authorization": f"Bearer {r.json()['access_token']}"}
+
+
+class Peer:
+ def __init__(self, hub_id: str):
+ self.hub_id = hub_id
+ self._sk = Ed25519PrivateKey.generate()
+ self.pk_pem = self._sk.public_key().public_bytes(
+ serialization.Encoding.PEM,
+ serialization.PublicFormat.SubjectPublicKeyInfo).decode()
+
+ def _sk_pem(self) -> bytes:
+ return self._sk.private_bytes(
+ serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8,
+ serialization.NoEncryption())
+
+ def envelope(self, jti: str | None = None) -> str:
+ now = int(time.time())
+ return jwt.encode(
+ {"iss": self.hub_id, "sub": self.hub_id,
+ "jti": jti or str(uuid.uuid4()), "iat": now, "exp": now + 300},
+ self._sk_pem(), algorithm="EdDSA")
+
+ def revocation(self, target: str, target_id: str) -> str:
+ return jwt.encode(
+ {"type": "revocation", "target": target, "target_id": target_id,
+ "iss": self.hub_id, "iat": int(time.time())},
+ self._sk_pem(), algorithm="EdDSA")
+
+ def header(self, **kw) -> dict:
+ return {"Authorization": f"Bearer {self.envelope(**kw)}"}
+
+
+async def _register_peer(client, admin, peer: Peer):
+ r = await client.post("/mhp/peers", headers=admin, json={
+ "hub_id": peer.hub_id, "hub_url": f"https://{peer.hub_id}",
+ "pk_hub_pem": peer.pk_pem})
+ assert r.status_code == 201, r.text
+
+
+# ── receive_directory ──────────────────────────────────────────────────────
+
+@pytest.mark.asyncio
+async def test_unknown_peer_is_refused(client):
+ stranger = Peer("nobody.example")
+ r = await client.post("/mhp/directory", headers=stranger.header(),
+ json={"hub_id": "nobody.example", "groups": []})
+ assert r.status_code == 401
+
+
+@pytest.mark.asyncio
+async def test_source_hub_is_the_signer_not_the_body(client):
+ admin = await _admin(client)
+ peer = Peer("peer-a.example")
+ await _register_peer(client, admin, peer)
+
+ r = await client.post("/mhp/directory", headers=peer.header(), json={
+ "hub_id": "peer-b.example", # claims to relay another hub
+ "groups": [{"id": "g-1", "name": "Shared", "join_policy": "open"}]})
+ assert r.status_code == 202
+
+ listing = (await client.get("/v1/groups")).json()["groups"]
+ row = next(g for g in listing if g["id"] == "g-1")
+ assert row["source"] == "peer-a.example" # the signer, not "peer-b.example"
+
+
+@pytest.mark.asyncio
+async def test_a_federated_id_cannot_shadow_a_local_group(client):
+ admin = await _admin(client)
+ peer = Peer("peer-a.example")
+ await _register_peer(client, admin, peer)
+
+ owner = await _admin(client, "owner")
+ r = await client.post("/v1/groups", headers=owner, json={
+ "name": "mine", "visibility": "public", "join_policy": "open"})
+ local_id = r.json()["group_id"]
+
+ r = await client.post("/mhp/directory", headers=peer.header(), json={
+ "hub_id": peer.hub_id,
+ "groups": [{"id": local_id, "name": "evil twin", "join_policy": "open"}]})
+ assert r.status_code == 202
+ assert r.json()["accepted"] == 0
+
+
+@pytest.mark.asyncio
+async def test_a_state_changing_token_cannot_be_replayed(client):
+ admin = await _admin(client)
+ peer = Peer("peer-a.example")
+ await _register_peer(client, admin, peer)
+
+ env = peer.envelope(jti="fixed-jti")
+ h = {"Authorization": f"Bearer {env}"}
+ body = {"hub_id": peer.hub_id,
+ "groups": [{"id": "g-9", "name": "Once", "join_policy": "open"}]}
+
+ assert (await client.post("/mhp/directory", headers=h, json=body)).status_code == 202
+ assert (await client.post("/mhp/directory", headers=h, json=body)).status_code == 401
+
+
+# ── receive_revocation ─────────────────────────────────────────────────────
+
+@pytest.mark.asyncio
+async def test_a_peer_may_withdraw_its_own_group(client):
+ admin = await _admin(client)
+ peer = Peer("peer-a.example")
+ await _register_peer(client, admin, peer)
+
+ await client.post("/mhp/directory", headers=peer.header(), json={
+ "hub_id": peer.hub_id,
+ "groups": [{"id": "g-77", "name": "Bye", "join_policy": "open"}]})
+ assert any(g["id"] == "g-77" for g in (await client.get("/v1/groups")).json()["groups"])
+
+ r = await client.post("/mhp/revoke", headers=peer.header(),
+ json={"token": peer.revocation("group", "g-77")})
+ assert r.status_code == 202 and r.json()["pruned"] == 1
+ assert not any(g["id"] == "g-77" for g in (await client.get("/v1/groups")).json()["groups"])
+
+
+@pytest.mark.asyncio
+async def test_a_peer_cannot_withdraw_another_hubs_group(client):
+ admin = await _admin(client)
+ a, b = Peer("peer-a.example"), Peer("peer-b.example")
+ await _register_peer(client, admin, a)
+ await _register_peer(client, admin, b)
+
+ await client.post("/mhp/directory", headers=a.header(), json={
+ "hub_id": a.hub_id,
+ "groups": [{"id": "g-a", "name": "A's", "join_policy": "open"}]})
+
+ # b signs a revocation for a's group and presents it under b's envelope.
+ r = await client.post("/mhp/revoke", headers=b.header(),
+ json={"token": b.revocation("group", "g-a")})
+ assert r.status_code == 202 and r.json()["pruned"] == 0
+ assert any(g["id"] == "g-a" for g in (await client.get("/v1/groups")).json()["groups"])
+
+
+@pytest.mark.asyncio
+async def test_federation_cannot_revoke_a_user(client):
+ admin = await _admin(client)
+ peer = Peer("peer-a.example")
+ await _register_peer(client, admin, peer)
+
+ r = await client.post("/mhp/revoke", headers=peer.header(),
+ json={"token": peer.revocation("user", "some-user-id")})
+ assert r.status_code == 202 and r.json()["pruned"] == 0
diff --git a/packages/meshbay-hub/tests/test_moderation.py b/packages/meshbay-hub/tests/test_moderation.py
index 93d23cd..6848929 100644
--- a/packages/meshbay-hub/tests/test_moderation.py
+++ b/packages/meshbay-hub/tests/test_moderation.py
@@ -1,34 +1,58 @@
-"""Tests for moderation — reports + blocklist."""
+"""Tests for moderation — reports + blocklist.
+
+Reporting requires a signed-in account (it used to be anonymous, which made it a
+network-wide censorship primitive), the auto-block threshold counts *distinct
+reporting accounts*, and the whole flow is refused when the hub has public groups
+switched off.
+"""
import pytest
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
-from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey
-from meshbay_common.crypto import pk_to_b64
from meshbay_hub.api.deps import set_admin_usernames
-
FAKE_HASH = "a" * 64 # valid blake3 hex
-@pytest.fixture
-async def auth_headers(client):
- sk_ed = Ed25519PrivateKey.generate()
- sk_x = X25519PrivateKey.generate()
+async def _register_and_login(client, username: str) -> dict:
await client.post("/v1/users/register", json={
- "username": "mod_admin", "email": "m@t.com", "password": "modpass99",
- "pk_user_ed25519": pk_to_b64(sk_ed.public_key()),
- "pk_user_x25519": pk_to_b64(sk_x.public_key()),
+ "username": username, "email": f"{username}@t.com",
+ "password": "reporter99pw",
})
r = await client.post("/v1/users/login",
- json={"username": "mod_admin", "password": "modpass99"})
- set_admin_usernames(["mod_admin"])
+ json={"username": username, "password": "reporter99pw"})
return {"Authorization": f"Bearer {r.json()['access_token']}"}
+@pytest.fixture
+async def reporter(client):
+ return await _register_and_login(client, "reporter_one")
+
+
+@pytest.fixture
+async def admin_headers(client):
+ headers = await _register_and_login(client, "mod_admin")
+ set_admin_usernames(["mod_admin"])
+ return headers
+
+
@pytest.mark.asyncio
-async def test_report_content_logged(client):
+async def test_report_requires_auth(client):
+ # No credentials at all — FastAPI rejects the missing header before the body.
r = await client.post("/v1/reports", json={
"content_hash": FAKE_HASH, "reason": "illegal"})
+ assert r.status_code in (401, 422)
+
+ # A bogus token is a clean 401.
+ r = await client.post("/v1/reports",
+ json={"content_hash": FAKE_HASH, "reason": "illegal"},
+ headers={"Authorization": "Bearer not-a-real-token"})
+ assert r.status_code == 401
+
+
+@pytest.mark.asyncio
+async def test_report_content_logged(client, reporter):
+ r = await client.post("/v1/reports",
+ json={"content_hash": FAKE_HASH, "reason": "illegal"},
+ headers=reporter)
assert r.status_code == 201
data = r.json()
assert data["report_count"] == 1
@@ -36,43 +60,68 @@ async def test_report_content_logged(client):
@pytest.mark.asyncio
-async def test_auto_block_on_threshold(client):
- """Second report triggers auto-block."""
- hash2 = "b" * 64
- await client.post("/v1/reports", json={"content_hash": hash2, "reason": "spam"})
- r = await client.post("/v1/reports", json={"content_hash": hash2, "reason": "spam"})
+async def test_same_reporter_cannot_walk_the_threshold(client, reporter):
+ h = "b" * 64
+ for _ in range(5):
+ r = await client.post("/v1/reports",
+ json={"content_hash": h, "reason": "spam"},
+ headers=reporter)
+ assert r.json()["report_count"] == 1
+ assert r.json()["status"] == "already_reported"
+
+ check = await client.get(f"/v1/blocklist/check?hash={h}")
+ assert check.json()["blocked"] is False
+
+
+@pytest.mark.asyncio
+async def test_auto_block_on_distinct_reporters(client):
+ h = "c" * 64
+ for i in range(3):
+ headers = await _register_and_login(client, f"rep_{i}")
+ r = await client.post("/v1/reports",
+ json={"content_hash": h, "reason": "illegal"},
+ headers=headers)
assert r.json()["status"] == "auto_blocked"
- assert r.json()["report_count"] == 2
+ assert r.json()["report_count"] == 3
+
+ check = await client.get(f"/v1/blocklist/check?hash={h}")
+ assert check.json()["blocked"] is True
@pytest.mark.asyncio
-async def test_blocklist_check(client):
- hash3 = "c" * 64
- # Not blocked yet
- r = await client.get(f"/v1/blocklist/check?hash={hash3}")
- assert r.json()["blocked"] is False
+async def test_reports_refused_when_public_groups_disabled(client, reporter, admin_headers):
+ await client.patch("/v1/admin/settings",
+ json={"allow_public_groups": False},
+ headers=admin_headers)
- # Report twice to auto-block
- await client.post("/v1/reports", json={"content_hash": hash3, "reason": "illegal"})
- await client.post("/v1/reports", json={"content_hash": hash3, "reason": "illegal"})
+ r = await client.post("/v1/reports",
+ json={"content_hash": "d" * 64, "reason": "illegal"},
+ headers=reporter)
+ assert r.status_code == 403
- r = await client.get(f"/v1/blocklist/check?hash={hash3}")
- assert r.json()["blocked"] is True
+
+@pytest.mark.asyncio
+async def test_invalid_hash_rejected(client, reporter):
+ r = await client.post("/v1/reports",
+ json={"content_hash": "not-a-valid-blake3-hash",
+ "reason": "test"},
+ headers=reporter)
+ assert r.status_code == 422
@pytest.mark.asyncio
-async def test_admin_add_remove_blocklist(client, auth_headers):
- hash4 = "d" * 64
+async def test_admin_add_remove_blocklist(client, admin_headers):
+ hash4 = "e" * 64
r = await client.post("/v1/admin/blocklist",
json={"content_hash": hash4, "reason": "csam"},
- headers=auth_headers)
+ headers=admin_headers)
assert r.status_code == 201
r = await client.get(f"/v1/blocklist/check?hash={hash4}")
assert r.json()["blocked"] is True
- r = await client.delete(f"/v1/admin/blocklist/{hash4}", headers=auth_headers)
+ r = await client.delete(f"/v1/admin/blocklist/{hash4}", headers=admin_headers)
assert r.status_code == 200
r = await client.get(f"/v1/blocklist/check?hash={hash4}")
@@ -80,18 +129,11 @@ async def test_admin_add_remove_blocklist(client, auth_headers):
@pytest.mark.asyncio
-async def test_invalid_hash_rejected(client):
- r = await client.post("/v1/reports", json={
- "content_hash": "not-a-valid-blake3-hash", "reason": "test"})
- assert r.status_code == 422
-
-
-@pytest.mark.asyncio
-async def test_full_blocklist(client, auth_headers):
- hash5 = "e" * 64
+async def test_full_blocklist(client, admin_headers):
+ hash5 = "f" * 64
await client.post("/v1/admin/blocklist",
json={"content_hash": hash5, "reason": "test"},
- headers=auth_headers)
+ headers=admin_headers)
r = await client.get("/v1/blocklist")
assert r.status_code == 200
assert hash5 in r.json()["hashes"]
diff --git a/packages/meshbay-hub/tests/test_node_auth.py b/packages/meshbay-hub/tests/test_node_auth.py
index 72ce412..a104a72 100644
--- a/packages/meshbay-hub/tests/test_node_auth.py
+++ b/packages/meshbay-hub/tests/test_node_auth.py
@@ -148,24 +148,35 @@ async def test_node_scope_blocks_group_create(client):
@pytest.mark.asyncio
-async def test_node_scope_blocks_add_member(client):
- sk_node, user_token = await _setup_node_user(client, "op1")
+async def test_node_token_may_add_a_member_to_its_own_operators_group(client):
+ """The node calls this after a CLI `member invite` so the group shows up in
+ the invitee's SPA (commit 0443cf8). A node-scoped token is accepted here —
+ the `group.admin_id == caller` check is the guard — but only for a group the
+ node's operator owns."""
+ sk_op, op_token = await _setup_node_user(client, "op1")
r = await client.post("/v1/groups", json={
"name": "mygroup", "visibility": "private", "join_policy": "invite",
- }, headers={"Authorization": f"Bearer {user_token}"})
- assert r.status_code == 201
+ }, headers={"Authorization": f"Bearer {op_token}"})
gid = r.json()["group_id"]
_, pk2 = _gen_ed25519()
_, px2 = _gen_x25519()
await _register(client, "member1", pk2, px2)
- r = await _node_auth(client, "op1", sk_node)
- node_token = r.json()["access_token"]
+ node_token = (await _node_auth(client, "op1", sk_op)).json()["access_token"]
r = await client.post(f"/v1/groups/{gid}/members/member1",
headers={"Authorization": f"Bearer {node_token}"})
+ assert r.status_code == 201
+
+ # …but not to a group it does not own.
+ sk_other, other_token = await _setup_node_user(client, "op2")
+ r = await client.post("/v1/groups", json={"name": "theirs", "visibility": "private"},
+ headers={"Authorization": f"Bearer {other_token}"})
+ other_gid = r.json()["group_id"]
+ r = await client.post(f"/v1/groups/{other_gid}/members/member1",
+ headers={"Authorization": f"Bearer {node_token}"})
assert r.status_code == 403
diff --git a/packages/meshbay-hub/tests/test_register_captcha.py b/packages/meshbay-hub/tests/test_register_captcha.py
new file mode 100644
index 0000000..befc1e2
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_register_captcha.py
@@ -0,0 +1,58 @@
+"""Registration CAPTCHA is enforced for every fresh account when configured.
+
+The gate used to be skipped whenever the request carried an `auth_key` — which
+every real client sends (the password split) — so it protected nobody and a bot
+skipped it by including the field. It now runs on `captcha.enabled` alone; the
+desktop client is Chromium and renders the same widget.
+"""
+
+import pytest
+
+
+@pytest.fixture
+def captcha_on(client, monkeypatch):
+ """Turn on a fake captcha: any config with both keys is `enabled`, and
+ verification succeeds only for the token 'good-token'."""
+ from meshbay_hub.api.users import _cfg
+ monkeypatch.setattr(_cfg.captcha, "site_key", "test-site")
+ monkeypatch.setattr(_cfg.captcha, "secret_key", "test-secret")
+
+ async def fake_verify(secret, token, remote_ip=None):
+ return token == "good-token"
+
+ monkeypatch.setattr("meshbay_hub.captcha.verify_captcha", fake_verify)
+
+
+def _body(**over):
+ b = {"username": "newbie", "email": "newbie@t.com", "auth_key": "a" * 44}
+ b.update(over)
+ return b
+
+
+@pytest.mark.asyncio
+async def test_missing_captcha_rejected_even_with_auth_key(client, captcha_on):
+ r = await client.post("/v1/users/register", json=_body())
+ assert r.status_code == 400
+ assert r.json()["detail"] == "captcha_required"
+
+
+@pytest.mark.asyncio
+async def test_bad_captcha_rejected(client, captcha_on):
+ r = await client.post("/v1/users/register",
+ json=_body(captcha_token="wrong"))
+ assert r.status_code == 400
+ assert r.json()["detail"] == "captcha_failed"
+
+
+@pytest.mark.asyncio
+async def test_good_captcha_accepted(client, captcha_on):
+ r = await client.post("/v1/users/register",
+ json=_body(captcha_token="good-token"))
+ assert r.status_code == 201
+
+
+@pytest.mark.asyncio
+async def test_no_captcha_configured_still_registers(client):
+ # Default test config has no captcha keys — registration proceeds without one.
+ r = await client.post("/v1/users/register", json=_body())
+ assert r.status_code == 201
diff --git a/packages/meshbay-hub/tests/test_security_headers.py b/packages/meshbay-hub/tests/test_security_headers.py
new file mode 100644
index 0000000..b4d7e6d
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_security_headers.py
@@ -0,0 +1,65 @@
+"""
+The hub sends a Content-Security-Policy and the other protective headers on
+every response — the SPA shell, its assets, and the API alike.
+
+Second-review L5 / third-review M5: previously there were none, so an injection
+that landed in the SPA (rendered third-party OG data, a federated group name,
+chat content) had nothing stopping it from loading more code or exfiltrating.
+"""
+
+import pytest
+from meshbay_hub.api.webapp import CSP
+
+
+def _directive(csp: str, name: str) -> str:
+ for part in csp.split(";"):
+ part = part.strip()
+ if part == name or part.startswith(name + " "):
+ return part
+ return ""
+
+
+@pytest.mark.asyncio
+async def test_the_spa_shell_carries_the_policy(client):
+ r = await client.get("/")
+ assert r.headers["content-security-policy"] == CSP
+ assert r.headers["x-content-type-options"] == "nosniff"
+ assert r.headers["x-frame-options"] == "DENY"
+ assert "referrer-policy" in r.headers
+
+
+@pytest.mark.asyncio
+async def test_the_api_carries_the_headers_too(client):
+ r = await client.get("/v1/health")
+ assert r.status_code == 200
+ assert "content-security-policy" in r.headers
+ assert r.headers["x-content-type-options"] == "nosniff"
+
+
+@pytest.mark.asyncio
+async def test_even_a_404_carries_the_headers(client):
+ # The middleware runs on every response, so a probe for a missing path
+ # cannot be framed or content-sniffed either.
+ r = await client.get("/no/such/path")
+ assert r.status_code == 404
+ assert r.headers["x-frame-options"] == "DENY"
+
+
+def test_the_policy_is_locked_down_where_it_matters():
+ assert "default-src 'none'" in CSP # covers object-src, etc.
+ assert _directive(CSP, "frame-ancestors") == "frame-ancestors 'none'"
+ assert _directive(CSP, "base-uri") == "base-uri 'none'"
+
+ script = _directive(CSP, "script-src")
+ # The hub's own origin must not be able to serve executable script (T3):
+ # 'self' and the wasm token are fine, a bare `https:` scheme is not.
+ assert "'self'" in script and "'wasm-unsafe-eval'" in script
+ assert "https:" not in script.split()
+
+
+def test_recaptcha_is_the_only_external_origin():
+ hosts = {"https://www.google.com", "https://www.gstatic.com"}
+ for part in CSP.split(";"):
+ for tok in part.strip().split()[1:]:
+ if tok.startswith(("http://", "https://")):
+ assert tok in hosts, f"unexpected external origin in CSP: {tok}"
diff --git a/packages/meshbay-node/src/meshbay_node/config.py b/packages/meshbay-node/src/meshbay_node/config.py
index b1ebd54..78ee873 100644
--- a/packages/meshbay-node/src/meshbay_node/config.py
+++ b/packages/meshbay-node/src/meshbay_node/config.py
@@ -36,7 +36,10 @@ url = "https://meshbay.org"
username = "myusername"
[node]
-quic_port = 19010 # QUIC (MNP) — LAN, port-forwarded, hub-less direct access
+# QUIC (MNP) direct path — LAN, port-forwarded, hub-less. Off by default: no
+# client speaks QUIC yet, so leaving it on only opens a UDP port.
+quic_enabled = false
+quic_port = 19010
ui_port = 18000 # local control API — JSON, 127.0.0.1 only, token-gated
# One-time codes. An invitation waits for someone to read their messages; an
@@ -129,6 +132,13 @@ class HubConfig:
@dataclass
class NodeConfig:
quic_port: int = 19010
+ # The QUIC MNP listener. Off by default: no shipping client speaks QUIC yet
+ # (the browser and the desktop client use WebRTC; the hub-less `group://`
+ # sidecar is unbuilt), so starting it only opens a UDP port with nothing to
+ # reach it. Turn on for LAN / port-forwarded / hub-less direct access once a
+ # client for it exists. `punch_nat()` is a direct-connection helper, not a
+ # NAT-traversal stack — a peer behind NAT still needs the port forwarded.
+ quic_enabled: bool = False
ui_port: int = 18000
# How long a one-time code stays usable. Invitations travel through a human
# conversation and are answered days later; operator pairing happens during
@@ -303,6 +313,7 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config:
# `port` (TCP+TLS) and `http_port` no longer exist — both listeners were removed
# in Phase 11.5 (findings C1, C6). Regenerate node.toml with `meshbay-node init`.
cfg.node.quic_port = nd.get("quic_port", cfg.node.quic_port)
+ cfg.node.quic_enabled = bool(nd.get("quic_enabled", cfg.node.quic_enabled))
cfg.node.ui_port = nd.get("ui_port", cfg.node.ui_port)
cfg.node.invite_ttl_hours = int(
nd.get("invite_ttl_hours", cfg.node.invite_ttl_hours))
@@ -371,6 +382,8 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config:
cfg.hub.username = user
if port := os.environ.get("MESHBAY_QUIC_PORT"):
cfg.node.quic_port = int(port)
+ if (qe := os.environ.get("MESHBAY_QUIC_ENABLED")) is not None:
+ cfg.node.quic_enabled = qe.strip().lower() in ("1", "true", "yes", "on")
if streams := os.environ.get("MESHBAY_MAX_CONCURRENT_STREAMS"):
cfg.node.max_concurrent_streams = _positive(
streams, cfg.node.max_concurrent_streams,
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index 056371a..8f9307c 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -537,7 +537,11 @@ class NodeDaemon:
log.warning("WebRTC not available (aiortc not installed)")
# 7. QUIC chunk server (LAN / port-forwarded / hub-less direct access)
- if QUIC_AVAILABLE:
+ #
+ # Off unless `[node] quic_enabled = true`: no shipping client speaks
+ # QUIC (browser and desktop use WebRTC; the `group://` sidecar is
+ # unbuilt), so starting it by default only exposes a UDP port.
+ if QUIC_AVAILABLE and self._config.node.quic_enabled:
self._quic_server = QuicChunkServer(
sk_node=keys.sk_ed25519,
hub_pk_pem=session.hub_pk_pem,
@@ -553,6 +557,8 @@ class NodeDaemon:
await self._quic_server.start()
log.info("QUIC server on port %d (%d groups)",
self._config.node.quic_port, len(groups_ctx))
+ elif QUIC_AVAILABLE:
+ 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):
@@ -1708,6 +1714,7 @@ def main() -> None:
f'username = "{username}"',
"",
"[node]",
+ "quic_enabled = false # QUIC direct path; no client uses it yet",
"quic_port = 19010",
"ui_port = 18000",
"",
diff --git a/packages/meshbay-node/src/meshbay_node/linkpreview.py b/packages/meshbay-node/src/meshbay_node/linkpreview.py
index 64067d1..b223aea 100644
--- a/packages/meshbay-node/src/meshbay_node/linkpreview.py
+++ b/packages/meshbay-node/src/meshbay_node/linkpreview.py
@@ -14,12 +14,15 @@ hub:
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
-credentials, and the 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. Residual: a DNS name that
-resolves clean here and to something internal microseconds later at connect
-time (rebinding) — narrow, and closed properly by pinning the checked IP,
-which is a follow-up.
+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`).
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).
@@ -42,9 +45,15 @@ _TIMEOUT = 5.0
_MAX_REDIRECTS = 3
_MAX_HTML_BYTES = 512 * 1024
_MAX_IMAGE_BYTES = 2 * 1024 * 1024
+_MAX_IMAGE_PIXELS = 40_000_000 # ~40 MP; an OG card image is a fraction of this
_IMAGE_MAX_DIM = 600
_UA = "MeshBayBot/1.0 (+https://meshbay.org; link preview)"
+# Ports a real OpenGraph-bearing page is served on. Everything else — SSH, mail,
+# databases, caches, search, admin panels — is refused, so a member cannot aim
+# the node at an arbitrary service even on a public host.
+_ALLOWED_PORTS = frozenset({80, 443, 8080, 8443})
+
class UnsafeURL(ValueError):
"""The URL points somewhere the node must not fetch from."""
@@ -75,6 +84,12 @@ def safe_url(url: str) -> str:
host = parts.hostname
if not host:
raise UnsafeURL("no host")
+ try:
+ port = parts.port
+ except ValueError:
+ 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.
@@ -136,12 +151,32 @@ 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."""
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)
if resp.is_redirect and "location" in resp.headers:
current = safe_url(urljoin(current, resp.headers["location"]))
continue
@@ -242,6 +277,10 @@ def _downscale(raw: bytes) -> bytes | None:
return None
try:
with Image.open(BytesIO(raw)) as im:
+ # The header is parsed but the pixels are not decoded yet — refuse a
+ # decompression bomb before convert()/thumbnail() allocate for it.
+ if im.width * im.height > _MAX_IMAGE_PIXELS:
+ return None
im = im.convert("RGB")
im.thumbnail((_IMAGE_MAX_DIM, _IMAGE_MAX_DIM))
out = BytesIO()
diff --git a/packages/meshbay-node/src/meshbay_node/transport/quic_server.py b/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
index ce6fe17..30c7daf 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
@@ -61,6 +61,15 @@ CHUNK_SIZE = 1024 * 1024
MAX_MSG = 64 * 1024 * 1024
ALPN = ["meshbay-mnp"]
+# ffmpeg is spawned per STREAM_SEGMENT request, and `_extract_segment` runs
+# `subprocess.run` synchronously — so without a bound, an authenticated peer can
+# both fork-bomb the node and block its event loop for up to 30 s per request
+# (finding M2c). Extraction now runs in a thread and passes through this
+# semaphore. Small on purpose: the QUIC path has no shipping client yet, this is
+# parity work with the WebRTC transcode cap.
+_MAX_CONCURRENT_SEGMENTS = 4
+_segment_sem = asyncio.Semaphore(_MAX_CONCURRENT_SEGMENTS)
+
class Denylist:
"""
@@ -204,6 +213,9 @@ class _MNPServerProtocol(QuicConnectionProtocol):
self._nonce_client: bytes = b""
self._gek_challenge: bytes | None = None
self._pending = None
+ # asyncio holds only a weak reference to a bare task, so a spawned
+ # handler still running can be collected mid-flight. Hold them.
+ self._tasks: set[asyncio.Task] = set()
def quic_event_received(self, event: QuicEvent) -> None:
if isinstance(event, StreamDataReceived):
@@ -232,7 +244,7 @@ class _MNPServerProtocol(QuicConnectionProtocol):
elif mtype == MNP.FILE_REQUEST:
self._do_file_request_sync(stream_id, msg)
elif mtype == MNP.STREAM_SEGMENT:
- self._do_stream_segment_sync(stream_id, msg)
+ self._spawn(self._do_stream_segment(stream_id, msg))
elif mtype == MNP.CHAT_MESSAGE:
self._do_chat_message_sync(stream_id, msg)
elif mtype == MNP.PING:
@@ -256,11 +268,10 @@ class _MNPServerProtocol(QuicConnectionProtocol):
checks could drift from the WebRTC path independently. All of that now
comes from meshbay_common.handshake, shared with WebRTC.
- NOT YET DONE — finding C6 remains open on this transport: there is still no
- GEK proof here, so a forged or stolen token reaches the node and can inject
- chat without holding the group key. The challenge/response and mutual node
- proof (quic_binding() is written and unit-tested for exactly this) are the
- remaining work in 11.5.4/5/6.
+ The GEK proof is enforced here too: `_do_handshake_response_sync` runs
+ the same challenge/response and mutual node proof, from the same shared
+ module, bound to the QUIC certificate hash. Finding C6 is closed on this
+ transport.
"""
try:
peer = authorize_token(
@@ -339,9 +350,7 @@ class _MNPServerProtocol(QuicConnectionProtocol):
self._user_id = peer.user_id
self._group_id = peer.group_id
- peers = self._ctx.get("_peers")
- if peers is not None:
- peers[self._user_id] = self
+ self._peer_registry()[self._user_id] = self
transcript = handshake_transcript(
ROLE_NODE, peer.group_id, self._nonce_client, self._gek_challenge, binding)
@@ -365,6 +374,20 @@ class _MNPServerProtocol(QuicConnectionProtocol):
return self._ctx["groups"][self._group_id]
return self._ctx
+ def _spawn(self, coro) -> None:
+ task = asyncio.ensure_future(coro)
+ self._tasks.add(task)
+ task.add_done_callback(self._tasks.discard)
+
+ def _peer_registry(self) -> dict:
+ """QUIC peers for THIS connection's group, keyed per group so a message
+ never crosses into another group on a multi-group node (findings M2b /
+ H1). Deliberately separate from the WebRTC registry that also lives in
+ the group context: the two transports' session objects have different
+ `_send` signatures, and cross-transport chat fan-out is not wired (no
+ QUIC client ships yet)."""
+ return self._group_ctx().setdefault("_quic_peers", {})
+
def _do_index_sync_sync(self, stream_id: int) -> None:
ctx = self._group_ctx()
wire = ctx["index"].serialize()
@@ -399,61 +422,83 @@ class _MNPServerProtocol(QuicConnectionProtocol):
)
self._send(stream_id, chunk_data)
- def _do_stream_segment_sync(self, stream_id: int, msg: dict) -> None:
- """Extract and serve one HLS segment via ffmpeg."""
- ctx = self._group_ctx()
- file_id = msg["file_id"]
- segment_index = msg["segment_index"]
- segment_duration = msg.get("segment_duration", 4)
+ async def _do_stream_segment(self, stream_id: int, msg: dict) -> None:
+ """
+ Extract and serve one segment via ffmpeg — off the event loop and behind
+ a concurrency bound, so one request can neither stall the whole node nor
+ fork-bomb it (finding M2c). The WebRTC path has had both since Phase 11.5.
+ """
+ try:
+ ctx = self._group_ctx()
+ file_id = msg["file_id"]
+ segment_index = msg["segment_index"]
+ segment_duration = msg.get("segment_duration", 4)
- entry = ctx["index"].get_entry(file_id)
- if not entry:
- self._send(stream_id, {"type": "error", "detail": "File not found"})
- return
+ entry = ctx["index"].get_entry(file_id)
+ if not entry:
+ self._send(stream_id, {"type": "error", "detail": "File not found"})
+ return
- file_path = entry_abs_path(ctx["roots"], entry)
- if not file_path.exists():
- self._send(stream_id, {"type": "error", "detail": "File not on disk"})
- return
+ file_path = entry_abs_path(ctx["roots"], entry)
+ if not file_path.exists():
+ self._send(stream_id, {"type": "error", "detail": "File not on disk"})
+ return
- start_time = segment_index * segment_duration
- segment_data = _extract_segment(file_path, start_time, segment_duration)
- if segment_data is None:
- self._send(stream_id, {"type": "error", "detail": "Segment extraction failed"})
- return
+ start_time = segment_index * segment_duration
+ loop = asyncio.get_event_loop()
+ async with _segment_sem:
+ segment_data = await loop.run_in_executor(
+ None, _extract_segment, file_path, start_time, segment_duration)
+ if segment_data is None:
+ self._send(stream_id, {"type": "error", "detail": "Segment extraction failed"})
+ return
- self._send(stream_id, {
- "type": MNP.STREAM_SEGMENT,
- "v": MNP_VERSION,
- "file_id": file_id,
- "segment_index": segment_index,
- "data_b64": base64.b64encode(segment_data).decode(),
- "size": len(segment_data),
- })
+ self._send(stream_id, {
+ "type": MNP.STREAM_SEGMENT,
+ "v": MNP_VERSION,
+ "file_id": file_id,
+ "segment_index": segment_index,
+ "data_b64": base64.b64encode(segment_data).decode(),
+ "size": len(segment_data),
+ })
+ except Exception as e:
+ log.error("stream_segment: %s", e)
+ self._send(stream_id, {"type": "error", "detail": "Segment extraction failed"})
def _do_chat_message_sync(self, stream_id: int, msg: dict) -> None:
- """Receive a chat message, store it, and broadcast to other connected peers."""
- chat_store = self._ctx.get("chat_store")
+ """
+ Store a chat message and broadcast it to the rest of THIS group.
+
+ `sender_id` is the authenticated session's, never the wire's — a peer
+ must not be able to post as someone else (NS6 / finding M2a). The store
+ and the peer set come from the group context, not a connection-global
+ one, so a message never crosses into another group on a multi-group node
+ (findings M2b / H1). The WebRTC path has done both since Phase 11.5.
+ """
+ gctx = self._group_ctx()
+ payload = msg.get("payload", b"")
+ if isinstance(payload, str):
+ payload = payload.encode()
+
+ chat_store = gctx.get("chat_store")
if chat_store:
- import asyncio
- asyncio.ensure_future(chat_store.save_message(
- sender_id=msg.get("sender_id", self._user_id),
+ self._spawn(chat_store.save_message(
+ sender_id=self._user_id,
iteration=msg.get("iteration", 0),
- payload=msg.get("payload", b"").encode() if isinstance(msg.get("payload"), str) else msg.get("payload", b""),
+ payload=payload,
thread_id=msg.get("thread_id"),
))
- peers = self._ctx.get("_peers", {})
broadcast = {
"type": MNP.CHAT_MESSAGE,
"v": MNP_VERSION,
- "sender_id": msg.get("sender_id", self._user_id),
+ "sender_id": self._user_id,
"iteration": msg.get("iteration", 0),
"payload": msg.get("payload", ""),
"thread_id": msg.get("thread_id"),
"group_id": self._group_id or "",
}
- for uid, proto in peers.items():
+ for uid, proto in list(self._peer_registry().items()):
if uid != self._user_id and proto is not self:
try:
proto._send(0, broadcast)
@@ -463,9 +508,10 @@ class _MNPServerProtocol(QuicConnectionProtocol):
self._send(stream_id, {"type": "ack", "v": MNP_VERSION})
def connection_lost(self, exc) -> None:
- peers = self._ctx.get("_peers")
- if peers and self._user_id:
- peers.pop(self._user_id, None)
+ if self._user_id:
+ self._peer_registry().pop(self._user_id, None)
+ for task in list(self._tasks):
+ task.cancel()
super().connection_lost(exc)
def _send(self, stream_id: int, obj: dict) -> None:
@@ -558,7 +604,7 @@ class QuicChunkServer:
self._ctx["groups"] = groups
self._denylist = denylist or Denylist()
self._ctx["denylist"] = self._denylist
- self._ctx["_peers"] = {}
+ # Peer sets are per group now — see _MNPServerProtocol._peer_registry().
self._host = host
self._port = port
self._cert_path = cert_path or Path.home() / ".config/meshbay/node_tls.crt"
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 003cd23..64ba75c 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -143,6 +143,16 @@ def _link_preview_cache_put(url: str, value: dict) -> None:
_link_preview_cache.pop(oldest, None)
_link_preview_cache[url] = (time.time(), value)
+
+# A member pasting a link is normal; a member — or a hub minting tokens for many
+# accounts — firing hundreds is amplification/DoS and a way to make the node
+# reach arbitrary hosts on demand (finding M3). Only a real outbound fetch is
+# counted (a cache hit costs nothing), and the ceilings are generous enough that
+# ordinary chat never meets them.
+_LINK_PREVIEW_RATE_WINDOW = 60.0
+_LINK_PREVIEW_RATE_PER_CONN = 15
+_LINK_PREVIEW_RATE_NODE = 60
+
# Upload limits (finding C5a). Uploads used to land directly in the shared root under
# a name the client chose, overwriting whatever was already there — which both violated
# node sovereignty and defeated the delete authorization (overwrite a file, become its
@@ -3588,6 +3598,27 @@ class WebRTCPeerSession:
],
})
+ def _link_preview_rate_ok(self) -> bool:
+ """
+ True when this preview fetch is within both the per-connection and the
+ node-wide window; records it when so, and both counts are trimmed to the
+ window on every call so the lists cannot grow without bound.
+ """
+ now = time.monotonic()
+ w = _LINK_PREVIEW_RATE_WINDOW
+ mine = [t for t in getattr(self, "_link_preview_hits", []) if now - t < w]
+ node = [t for t in self._ctx.get("link_preview_hits", []) if now - t < w]
+ if (len(mine) >= _LINK_PREVIEW_RATE_PER_CONN
+ or len(node) >= _LINK_PREVIEW_RATE_NODE):
+ self._link_preview_hits = mine
+ self._ctx["link_preview_hits"] = node
+ return False
+ mine.append(now)
+ node.append(now)
+ self._link_preview_hits = mine
+ self._ctx["link_preview_hits"] = node
+ return True
+
async def _do_link_preview_request(self, msg: dict) -> None:
"""
Unfurl a URL a member pasted into chat (draft-v6 §2.7 enrichment rule:
@@ -3608,6 +3639,15 @@ class WebRTCPeerSession:
self._send({**cached, "type": MNP.LINK_PREVIEW_RESP, "v": MNP_VERSION})
return
+ if not self._link_preview_rate_ok():
+ # Same shape as any other miss — the client shows the bare link. A
+ # rate-limited result is not cached, so it is retried once the
+ # window clears rather than pinned as "no preview".
+ log.debug("link_preview_req: rate-limited (peer=%s)", self._peer_id)
+ self._send({"type": MNP.LINK_PREVIEW_RESP, "v": MNP_VERSION,
+ "url": key, "ok": False})
+ return
+
resp: dict = {"type": MNP.LINK_PREVIEW_RESP, "v": MNP_VERSION,
"url": key, "ok": False}
try:
diff --git a/packages/meshbay-node/tests/test_link_preview_request.py b/packages/meshbay-node/tests/test_link_preview_request.py
index d975f91..fe7dec6 100644
--- a/packages/meshbay-node/tests/test_link_preview_request.py
+++ b/packages/meshbay-node/tests/test_link_preview_request.py
@@ -34,9 +34,10 @@ def _clear_cache():
webrtc_server._link_preview_cache.clear()
-def _session(media_cache):
+def _session(media_cache, ctx=None):
s = WebRTCPeerSession.__new__(WebRTCPeerSession)
- s._ctx = {"media_cache": media_cache}
+ s._ctx = ctx if ctx is not None else {"media_cache": media_cache}
+ s._peer_id = "t"
s.sent = []
s._send = s.sent.append
return s
@@ -79,6 +80,51 @@ async def test_unfurlable_failure_is_ok_false(media_cache, monkeypatch):
assert "image_thumb_hash" not in resp
+async def test_rate_limit_per_connection(media_cache, monkeypatch):
+ """A member firing many previews is bounded; over the ceiling the reply is
+ a plain `ok: false` (bare link) and no outbound fetch is made."""
+ monkeypatch.setattr(webrtc_server, "_LINK_PREVIEW_RATE_PER_CONN", 3)
+ calls = {"n": 0}
+
+ async def counting_preview(url, **k):
+ calls["n"] += 1
+ return {"url": url, "title": "x", "description": None,
+ "site_name": None, "image_url": None}
+ monkeypatch.setattr(linkpreview, "fetch_preview", counting_preview)
+
+ s = _session(media_cache)
+ for i in range(3):
+ await s._do_link_preview_request({"url": f"https://example.com/{i}"})
+ assert calls["n"] == 3
+ assert all(r["ok"] for r in s.sent)
+
+ await s._do_link_preview_request({"url": "https://example.com/over"})
+ assert calls["n"] == 3 # not fetched
+ assert s.sent[-1]["ok"] is False
+
+
+async def test_rate_limit_is_node_wide(media_cache, monkeypatch):
+ """Two connections share the node-wide ceiling."""
+ monkeypatch.setattr(webrtc_server, "_LINK_PREVIEW_RATE_PER_CONN", 100)
+ monkeypatch.setattr(webrtc_server, "_LINK_PREVIEW_RATE_NODE", 2)
+ calls = {"n": 0}
+
+ async def counting_preview(url, **k):
+ calls["n"] += 1
+ return {"url": url, "title": "x", "description": None,
+ "site_name": None, "image_url": None}
+ monkeypatch.setattr(linkpreview, "fetch_preview", counting_preview)
+
+ ctx = {"media_cache": media_cache}
+ a, b = _session(media_cache, ctx), _session(media_cache, ctx)
+ await a._do_link_preview_request({"url": "https://example.com/a"})
+ await b._do_link_preview_request({"url": "https://example.com/b"})
+ await b._do_link_preview_request({"url": "https://example.com/c"})
+
+ assert calls["n"] == 2
+ assert b.sent[-1]["ok"] is False
+
+
async def test_second_request_for_the_same_url_is_served_from_cache(media_cache, monkeypatch):
calls = {"n": 0}
diff --git a/packages/meshbay-node/tests/test_linkpreview.py b/packages/meshbay-node/tests/test_linkpreview.py
index 0dd951a..a14b173 100644
--- a/packages/meshbay-node/tests/test_linkpreview.py
+++ b/packages/meshbay-node/tests/test_linkpreview.py
@@ -48,6 +48,30 @@ def test_safe_url_refuses(url):
safe_url(url)
+@pytest.mark.parametrize("url", [
+ "http://example.com:22/x", # SSH
+ "http://example.com:3306/x", # MySQL
+ "http://example.com:6379/x", # Redis
+ "http://example.com:9200/x", # Elasticsearch
+ "http://example.com:5000/x", # a common internal admin port
+])
+def test_safe_url_refuses_non_web_ports(url, resolves_public):
+ with pytest.raises(UnsafeURL):
+ safe_url(url)
+
+
+@pytest.mark.parametrize("url", [
+ "http://example.com/x", # implicit 80
+ "https://example.com/x", # implicit 443
+ "http://example.com:80/x",
+ "https://example.com:443/x",
+ "http://example.com:8080/x",
+ "https://example.com:8443/x",
+])
+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"
@@ -153,3 +177,19 @@ async def test_fetch_image_downscales(resolves_public):
assert jpeg and jpeg[:2] == b"\xff\xd8" # JPEG SOI
with Image.open(BytesIO(jpeg)) as im:
assert max(im.size) <= linkpreview._IMAGE_MAX_DIM
+
+
+async def test_fetch_image_refuses_a_decompression_bomb(resolves_public, monkeypatch):
+ from io import BytesIO
+
+ from PIL import Image
+ # A tiny file that reports enormous dimensions from its header alone.
+ monkeypatch.setattr(linkpreview, "_MAX_IMAGE_PIXELS", 1_000_000)
+ buf = BytesIO()
+ Image.new("RGB", (2000, 2000), (0, 0, 0)).save(buf, format="PNG") # 4 MP > cap
+ bomb = buf.getvalue()
+
+ def handler(request):
+ return httpx.Response(200, headers={"content-type": "image/png"}, content=bomb)
+ async with _client(handler) as c:
+ assert await linkpreview.fetch_image("https://example.com/x.png", client=c) is None
diff --git a/packages/meshbay-node/tests/test_quic_enabled.py b/packages/meshbay-node/tests/test_quic_enabled.py
new file mode 100644
index 0000000..c0c1568
--- /dev/null
+++ b/packages/meshbay-node/tests/test_quic_enabled.py
@@ -0,0 +1,53 @@
+"""
+The QUIC MNP listener is off unless the operator turns it on.
+
+No shipping client speaks QUIC (browser and desktop use WebRTC; the hub-less
+`group://` sidecar is unbuilt), so a node that started it by default would only
+be exposing a UDP port. `daemon.py` gates `QuicChunkServer` on
+`self._config.node.quic_enabled`; these follow the value down the config path.
+"""
+
+import textwrap
+from pathlib import Path
+
+from meshbay_node.config import load_config
+
+
+def _cfg(tmp_path: Path, body: str):
+ p = tmp_path / "node.toml"
+ p.write_text(textwrap.dedent(body))
+ return load_config(p)
+
+
+def test_off_by_default(tmp_path):
+ cfg = _cfg(tmp_path, """
+ [node]
+ quic_port = 19010
+ """)
+ assert cfg.node.quic_enabled is False
+
+
+def test_the_operator_turns_it_on(tmp_path):
+ cfg = _cfg(tmp_path, """
+ [node]
+ quic_enabled = true
+ """)
+ assert cfg.node.quic_enabled is True
+
+
+def test_the_environment_can_force_it_on(tmp_path, monkeypatch):
+ monkeypatch.setenv("MESHBAY_QUIC_ENABLED", "1")
+ cfg = _cfg(tmp_path, """
+ [node]
+ quic_enabled = false
+ """)
+ assert cfg.node.quic_enabled is True
+
+
+def test_the_environment_can_force_it_off(tmp_path, monkeypatch):
+ monkeypatch.setenv("MESHBAY_QUIC_ENABLED", "false")
+ cfg = _cfg(tmp_path, """
+ [node]
+ quic_enabled = true
+ """)
+ assert cfg.node.quic_enabled is False