aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/groups.py25
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py53
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/signaling.py34
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/users.py112
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)