diff options
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/groups.py | 25 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/revocation.py | 53 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/signaling.py | 34 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/users.py | 112 |
4 files changed, 177 insertions, 47 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py index 10049b2..fc117bf 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/groups.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py @@ -1,6 +1,7 @@ """Group endpoints — /v1/groups/*""" import re +import unicodedata from datetime import UTC, datetime from fastapi import APIRouter, Depends, HTTPException, Query, Request @@ -347,6 +348,25 @@ async def join_group( "owner_username": owner} +# The column's width. A longer name was a database error on PostgreSQL (a 500) +# and silently truncated on SQLite. +MAX_GROUP_NAME = 128 +# Line breaks and other C0/C1 controls, and the bidirectional overrides that +# make a name display as something other than what it is. Joiners stay: an +# emoji family is a ZWJ sequence. +_BIDI_CONTROLS = frozenset("\u202a\u202b\u202c\u202d\u202e\u2066\u2067\u2068\u2069") + + +def _group_name_problem(name: str) -> str | None: + if not name: + return "A group needs a name." + if len(name) > MAX_GROUP_NAME: + return f"A group name is at most {MAX_GROUP_NAME} characters." + if any(unicodedata.category(c) == "Cc" or c in _BIDI_CONTROLS for c in name): + return "A group name cannot contain control characters." + return None + + class GroupCreateRequest(BaseModel): name: str visibility: str = "private" # public|private @@ -426,8 +446,9 @@ async def create_group( "anyone to be able to join.") name = body.name.strip() - if not name: - raise HTTPException(status_code=422, detail="A group needs a name.") + problem = _group_name_problem(name) + if problem: + raise HTTPException(status_code=422, detail=problem) # One name per owner, case-insensitively. Two *different* owners may each # have a "photos" — that is why the check is scoped to `admin_id` and why # the group's real identity stays its UUID. The DB has a unique index too diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py index 6a6baa7..5a33d77 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py @@ -87,6 +87,16 @@ NOTIFY_WINDOW_SECONDS = 60 _notify_window: dict[str, tuple[float, int]] = {} # node_id → (window start, count) +# `update_groups` re-reads the node's groups from the database. A node sends one +# when its configuration is reloaded — an operator attaching a group, an owner's +# approval arriving — so ten a minute is far past real use, and the same budget +# rule as chat_notify keeps one node from spending the hub's database for others. +UPDATE_GROUPS_BURST = 10 +_update_window: dict[str, tuple[float, int]] = {} +# The groups one node may claim in one message. An operator with fifty groups +# is a large one. +MAX_CLAIMED_GROUPS = 1000 + def forget_node(node_id: str) -> None: """Drop everything a disconnected node's socket owned. @@ -106,25 +116,35 @@ def forget_node(node_id: str) -> None: _node_users.pop(node_id, None) -def _notify_budget(node_id: str) -> bool: - """True if this node may send one more chat_notify now.""" +def _spend(window: dict[str, tuple[float, int]], node_id: str, burst: int) -> bool: + """True if this node may send one more message of a budgeted kind now.""" now = time.monotonic() - if len(_notify_window) > 1000: + if len(window) > 1000: # Swept here rather than on disconnect, which would let a node refill # its budget by reconnecting — the same token stays valid for an hour. - for nid, (started, _) in list(_notify_window.items()): + for nid, (started, _) in list(window.items()): if now - started >= NOTIFY_WINDOW_SECONDS: - _notify_window.pop(nid, None) - start, count = _notify_window.get(node_id, (now, 0)) + window.pop(nid, None) + start, count = window.get(node_id, (now, 0)) if now - start >= NOTIFY_WINDOW_SECONDS: start, count = now, 0 - if count >= NOTIFY_BURST: - _notify_window[node_id] = (start, count) + if count >= burst: + window[node_id] = (start, count) return False - _notify_window[node_id] = (start, count + 1) + window[node_id] = (start, count + 1) return True +def _notify_budget(node_id: str) -> bool: + """True if this node may send one more chat_notify now.""" + return _spend(_notify_window, node_id, NOTIFY_BURST) + + +def _update_budget(node_id: str) -> bool: + """True if this node may send one more update_groups now.""" + return _spend(_update_window, node_id, UPDATE_GROUPS_BURST) + + async def _mark_hosted(group_ids: list[str]) -> None: """Stamp the first time a node announced it hosts each of these groups. @@ -461,8 +481,8 @@ async def node_websocket(ws: WebSocket): await _reject(ws, "Node already connected", 4009) return - resolved_id, result = await _authorize_node_ws( - msg["token"], claimed_id, msg.get("group_ids")) + claimed = [str(g) for g in (msg.get("group_ids") or [])][:MAX_CLAIMED_GROUPS] + resolved_id, result = await _authorize_node_ws(msg["token"], claimed_id, claimed) if resolved_id is None: await _reject(ws, result, 4003) return @@ -472,7 +492,7 @@ async def node_websocket(ws: WebSocket): node_id = resolved_id _connected_nodes[node_id] = ws _node_groups[node_id] = group_ids - _node_claims[node_id] = list(msg.get("group_ids") or []) + _node_claims[node_id] = claimed _node_users[node_id] = user_id await _mark_hosted(group_ids) log.info("Node WS connected: %s (user=%s, groups=%d)", @@ -493,13 +513,16 @@ async def node_websocket(ws: WebSocket): from meshbay_hub.api.signaling import handle_webrtc_answer handle_webrtc_answer(msg, node_id) elif msg.get("type") == "update_groups": + if not _update_budget(node_id): + log.warning("Node %s exceeded its update_groups rate", node_id[:8]) + continue # Through the same gate as the registration above. This used to # assign the message's list verbatim, so the ceiling that makes # C2 hold at authentication could be stepped over one message # later: a node had only to reload to claim any group on the hub. - _node_claims[node_id] = list(msg.get("group_ids") or []) - new_gids = await resolve_node_groups( - node_id, user_id, msg.get("group_ids")) + claimed = [str(g) for g in (msg.get("group_ids") or [])][:MAX_CLAIMED_GROUPS] + _node_claims[node_id] = claimed + new_gids = await resolve_node_groups(node_id, user_id, claimed) _node_groups[node_id] = new_gids await _mark_hosted(new_gids) log.info("Node %s updated groups: %d", node_id[:8], len(new_gids)) diff --git a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py index 6c9699b..56e0e4b 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py @@ -20,7 +20,7 @@ import time import uuid from fastapi import APIRouter, Depends, HTTPException, Request -from pydantic import BaseModel +from pydantic import BaseModel, Field, field_validator from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession @@ -41,9 +41,23 @@ _webrtc_answers: dict[str, asyncio.Future] = {} _answer_owner: dict[str, str] = {} +# A browser offers a handful of candidates — a host and a reflexive one per +# interface — and embeds them in the SDP anyway. The list is relayed to the node +# as it came, so it is bounded like the SDP beside it. +MAX_ICE_CANDIDATES = 64 +MAX_ICE_BYTES = 32 * 1024 + + class WebRTCOfferRequest(BaseModel): sdp: str - ice_candidates: list[dict] = [] + ice_candidates: list[dict] = Field(default_factory=list, max_length=MAX_ICE_CANDIDATES) + + @field_validator("ice_candidates") + @classmethod + def _bounded(cls, v: list[dict]) -> list[dict]: + if len(json.dumps(v)) > MAX_ICE_BYTES: + raise ValueError("ICE candidates too large") + return v class WebRTCOfferResponse(BaseModel): @@ -176,14 +190,6 @@ async def webrtc_offer( if len(body.sdp) > MAX_SDP_BYTES: raise HTTPException(status_code=413, detail="SDP too large") - # Logged here because this is the moment a browser starts a peer connection, - # and the address it starts it from is this one — the hub's own view of the - # TCP connection. Whatever address the peers then discover through STUN is - # theirs to negotiate and is not what a log should record. - db.add(IPLog(user_id=current_user.id, event="webrtc_offer", - ip_address=client_ip(request), detail=node_id[:8])) - await db.commit() - ws = _connected_nodes.get(node_id) if not ws: raise HTTPException(status_code=404, detail="Node not connected") @@ -205,6 +211,14 @@ async def webrtc_offer( raise HTTPException(status_code=429, detail="Too many connections to this node", headers={"Retry-After": str(max(1, math.ceil(wait)))}) + # Logged once the offer is going to a node, not before: the address a peer + # connection starts from is the hub's own view of this TCP connection, and an + # IP log row is kept a year — written before the checks above, any account + # could add rows for any string it named as a node. + db.add(IPLog(user_id=current_user.id, event="webrtc_offer", + ip_address=client_ip(request), detail=node_id[:8])) + await db.commit() + peer_id = str(uuid.uuid4()) answer_future: asyncio.Future = asyncio.get_event_loop().create_future() _webrtc_answers[peer_id] = answer_future diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py index 8e780df..9326bfb 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/users.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py @@ -1,6 +1,7 @@ """User endpoints — /v1/users/*""" import base64 +import hashlib import logging import re import secrets @@ -42,6 +43,7 @@ from meshbay_hub.db.models import ( GroupInviteLink, GroupMember, IPLog, + KnownBrowser, Node, Notification, RefreshToken, @@ -161,6 +163,7 @@ class LoginRequest(BaseModel): username: str password: str | None = None # legacy (raw password) for migration auth_key: str | None = None # PBKDF2-derived auth key (new scheme) + known_browser: str | None = None # from an earlier sign-in on this browser class RefreshRequest(BaseModel): @@ -182,12 +185,18 @@ async def register( ): eh = hash_email_blind(body.email) - existing = await db.execute( - select(User).where(User.username == body.username)) - found = existing.scalar_one_or_none() + # Unique regardless of case: invitations and member management name people + # by username, and "Alice" beside "alice" is one person to whoever reads it. + # Accounts that already differ only by case (made before this) keep their + # names; the exact match is the one a retry means. + same = (await db.execute( + select(User).where(func.lower(User.username) == body.username.lower()) + )).scalars().all() + found = next((u for u in same if u.username == body.username), same[0] if same else None) if found: - if found.status == "pending" and found.email_hash == eh: + if (found.username == body.username and found.status == "pending" + and found.email_hash == eh): # Same person retrying before validation — resend a code. # No captcha: the initial registration already passed it. # @@ -351,6 +360,58 @@ async def _take_login_attempt(db: AsyncSession, username: str) -> None: headers={"Retry-After": str(retry_after)}) +def _session_counter(user: User) -> str: + """The failure counter for a passphrase re-checked inside an open session. + + Its own, not the sign-in one: a stranger who keeps a name locked at sign-in + must not also stop its owner changing their passphrase, deleting their + account or registering a device from a session they already hold. + """ + return f"\x00session:{user.id}" + + +async def _browser_counter(db: AsyncSession, username: str, + token: str | None) -> tuple[str, "KnownBrowser | None"]: + """The failure counter for a sign-in, and the known browser behind it if any. + + A browser that signed in to this account before presents its token and is + counted on its own: the username's counter, which anyone can spend, then + locks only browsers this account has never used. A token for another + account, or none, is the username's counter — the answer is the same either + way, so it says nothing about the account (M1). + """ + if token: + row = (await db.execute( + select(KnownBrowser).join(User, User.id == KnownBrowser.user_id) + .where(KnownBrowser.token_hash == _browser_hash(token), + User.username == username))).scalar_one_or_none() + if row is not None: + return f"{username}\x00browser:{row.id}", row + return username, None + + +def _browser_hash(token: str) -> str: + return hashlib.sha256(f"meshbay:known_browser:{token}".encode()).hexdigest() + + +# How many browsers one account is remembered on. The oldest goes first; a +# browser forgotten here is only an unknown one again. +MAX_KNOWN_BROWSERS = 20 + + +async def _remember_browser(db: AsyncSession, user: User) -> str: + """A new known-browser token for `user`. The caller commits.""" + raw = secrets.token_urlsafe(32) + rows = (await db.execute( + select(KnownBrowser.id).where(KnownBrowser.user_id == user.id) + .order_by(KnownBrowser.last_used_at.desc()))).scalars().all() + stale = rows[MAX_KNOWN_BROWSERS - 1:] + if stale: + await db.execute(delete(KnownBrowser).where(KnownBrowser.id.in_(stale))) + db.add(KnownBrowser(user_id=user.id, token_hash=_browser_hash(raw))) + return raw + + async def _prove_passphrase(db: AsyncSession, user: User, auth_key: str) -> None: """Refuse with 403 unless `auth_key` is this account's, spending an attempt. @@ -358,18 +419,18 @@ async def _prove_passphrase(db: AsyncSession, user: User, auth_key: str) -> None refreshed one, or one lifted from a page, and what it would buy here outlives the session or reopens the offline search the pepper exists to prevent. """ - await _take_login_attempt(db, user.username) + await _take_login_attempt(db, _session_counter(user)) if not await verify_password_off_loop(auth_key, user.pw_hash, user.pw_salt, user.pw_version): raise HTTPException(status_code=403, detail="Passphrase does not match") - await login_throttle.clear(db, user.username) + await login_throttle.clear(db, _session_counter(user)) async def _login_failed(db: AsyncSession, username: str, ip: str, - user_id: str | None = None) -> None: + user_id: str | None = None, counter: str | None = None) -> None: """Record a wrong passphrase and answer 401. Always raises.""" db.add(IPLog(user_id=user_id, event="login_fail", ip_address=ip, detail=username)) - if await login_throttle.is_now_locked(db, username): + if await login_throttle.is_now_locked(db, counter or username): # Once, on the failure that spent the last attempt — so the logs tab # shows when a name was locked, not every refusal after it. db.add(IPLog(user_id=user_id, event="login_locked", ip_address=ip, @@ -431,32 +492,33 @@ async def login( # Before the account is even looked up: an unknown name spends attempts and # locks exactly like a real one, so neither answer tells them apart (M1). - await _take_login_attempt(db, body.username) + counter, browser = await _browser_counter(db, body.username, body.known_browser) + await _take_login_attempt(db, counter) result = await db.execute( select(User).where(User.username == body.username)) user = result.scalar_one_or_none() if not user: - await _login_failed(db, body.username, ip) + await _login_failed(db, body.username, ip, counter=counter) if user.pw_version >= 3: # New scheme: verify auth_key if not body.auth_key or not await verify_password_off_loop( body.auth_key, user.pw_hash, user.pw_salt, version=user.pw_version ): - await _login_failed(db, body.username, ip, user.id) + await _login_failed(db, body.username, ip, user.id, counter) else: # Legacy scheme: need raw password if not body.password: # Nothing was checked, so nothing was guessed. - await login_throttle.release(db, body.username) + await login_throttle.release(db, counter) await db.commit() raise HTTPException(status_code=401, detail="auth_upgrade_required") if not await verify_password_off_loop( body.password, user.pw_hash, user.pw_salt, version=user.pw_version ): - await _login_failed(db, body.username, ip, user.id) + await _login_failed(db, body.username, ip, user.id, counter) # Migrate to new scheme if auth_key provided alongside password if body.auth_key: new_hash, new_salt = await hash_password_off_loop(body.auth_key) @@ -471,6 +533,7 @@ async def login( user.pw_version = 2 # The passphrase was right, whatever the account's status turns out to be. + await login_throttle.clear(db, counter) await login_throttle.clear(db, body.username) if user.status != "active": @@ -501,6 +564,11 @@ async def login( )) db.add(IPLog(user_id=user.id, event="login", ip_address=ip)) pepper = _bundle_pepper(user) + if browser is not None: + browser.last_used_at = datetime.now(UTC) + known = {} + else: + known = {"known_browser": await _remember_browser(db, user)} await db.commit() return { @@ -509,6 +577,7 @@ async def login( "token_type": "bearer", "expires_in": _ttl(), **pepper, + **known, } @@ -787,7 +856,8 @@ async def get_current_user_info( # A passphrase change re-wraps every node's bundle *before* the hub # accepts the new passphrase, and must not start while the hub would # then refuse it. - "passphrase_locked_for": await login_throttle.locked_for(db, current_user.username), + "passphrase_locked_for": await login_throttle.locked_for( + db, _session_counter(current_user)), } @@ -849,13 +919,13 @@ async def update_profile( raise HTTPException( status_code=403, detail="Changing your e-mail requires your passphrase.") - await _take_login_attempt(db, current_user.username) + await _take_login_attempt(db, _session_counter(current_user)) if not await verify_password_off_loop( body.auth_key, current_user.pw_hash, current_user.pw_salt, current_user.pw_version): raise HTTPException(status_code=403, detail="Passphrase does not match") - await login_throttle.clear(db, current_user.username) + await login_throttle.clear(db, _session_counter(current_user)) # How often one account may point the hub at a *different* address. # Long, because this is the only path where a signed-in account chooses @@ -1074,12 +1144,12 @@ async def change_password( current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): - await _take_login_attempt(db, current_user.username) + await _take_login_attempt(db, _session_counter(current_user)) if not await verify_password_off_loop(body.old_auth_key, current_user.pw_hash, current_user.pw_salt, current_user.pw_version): raise HTTPException(status_code=403, detail="Current passphrase does not match") - await login_throttle.clear(db, current_user.username) + await login_throttle.clear(db, _session_counter(current_user)) if body.new_auth_key == body.old_auth_key: raise HTTPException(status_code=400, detail="New passphrase must differ from the current one") @@ -1279,6 +1349,7 @@ async def password_reset( update(RefreshToken).where(RefreshToken.user_id == user.id) .values(revoked=True)) await db.execute(delete(UserDevice).where(UserDevice.user_id == user.id)) + await db.execute(delete(KnownBrowser).where(KnownBrowser.user_id == user.id)) # A code sent to the address on file is a stronger proof than a passphrase, # and it is the way out of a lockout somebody else caused. await login_throttle.clear(db, user.username) @@ -1493,6 +1564,7 @@ async def erase_account(db: AsyncSession, user: User, owned_groups: str = "refus GroupHost.node_id.in_(select(Node.id).where(Node.user_id == user.id)))) await db.execute(delete(Node).where(Node.user_id == user.id)) await db.execute(delete(UserDevice).where(UserDevice.user_id == user.id)) + await db.execute(delete(KnownBrowser).where(KnownBrowser.user_id == user.id)) await db.execute(delete(EmailVerification).where(EmailVerification.user_id == user.id)) # Links this account issued for a group it no longer owns; the ones for its # own groups went with them above. A used link keeps pointing at the @@ -1538,11 +1610,11 @@ async def delete_own_account( borrowed laptop or a session left open. Same value as at sign-in, so the hub still never sees the passphrase itself. """ - await _take_login_attempt(db, current_user.username) + await _take_login_attempt(db, _session_counter(current_user)) if not await verify_password_off_loop(body.auth_key, current_user.pw_hash, current_user.pw_salt, current_user.pw_version): raise HTTPException(status_code=403, detail="Passphrase does not match") - await login_throttle.clear(db, current_user.username) + await login_throttle.clear(db, _session_counter(current_user)) return await erase_account(db, current_user) |