diff options
Diffstat (limited to 'packages')
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 |