diff options
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/admin.py | 25 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/deps.py | 11 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/groups.py | 112 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/moderation.py | 318 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/relay.py | 166 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/revocation.py | 29 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/signaling.py | 2 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/users.py | 7 |
8 files changed, 290 insertions, 380 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/admin.py b/packages/meshbay-hub/src/meshbay_hub/api/admin.py index 381378c..dbda197 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/admin.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/admin.py @@ -43,6 +43,8 @@ class SettingsPatchRequest(BaseModel): login: dict[str, int] | None = None # Session lifetime, in hours, each optional. session: dict[str, int] | None = None + # Content reports: who may report, how often, what a report leads to. + reports: dict[str, int] | None = None # ── Instance settings ──────────────────────────────────────────────────────── @@ -63,6 +65,9 @@ async def _settings_payload(db: AsyncSession) -> dict: "session": await hub_settings.session_limits(db), "session_defaults": dict(hub_settings.SESSION_DEFAULTS), "session_bounds": {k: list(v) for k, v in hub_settings.SESSION_BOUNDS.items()}, + "reports": await hub_settings.report_limits(db), + "reports_defaults": dict(hub_settings.REPORT_DEFAULTS), + "reports_bounds": {k: list(v) for k, v in hub_settings.REPORT_BOUNDS.items()}, } @@ -162,6 +167,26 @@ async def admin_patch_settings( )) await db.commit() + if body.reports: + unknown = sorted(set(body.reports) - set(hub_settings.REPORT_KEYS)) + if unknown: + raise HTTPException( + status_code=422, detail=f"Unknown report setting(s): {unknown}") + changed = [] + for key, value in body.reports.items(): + clamped = hub_settings.clamp_report_value(key, value) + await hub_settings.set_raw(db, f"reports.{key}", str(clamped)) + changed.append(f"{key}={clamped}") + log.info("Report policy changed by %s: %s", + current_user.username, ", ".join(changed)) + db.add(IPLog( + user_id=current_user.id, + event="admin_reports_update", + ip_address="admin", + detail=", ".join(changed)[:255], + )) + await db.commit() + return await _settings_payload(db) diff --git a/packages/meshbay-hub/src/meshbay_hub/api/deps.py b/packages/meshbay-hub/src/meshbay_hub/api/deps.py index 42f4101..2fd8b99 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/deps.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/deps.py @@ -90,6 +90,17 @@ async def require_user_scope( return current_user +async def require_node_scope( + payload: dict = Depends(_decode_token), + current_user: User = Depends(get_current_user), +) -> User: + """Only a node daemon's own token — for what a node fetches on its own behalf.""" + if payload.get("scope") != "node": + raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, + detail="Node token required") + 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 diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py index 3a11345..df2f336 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/groups.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py @@ -20,16 +20,11 @@ from meshbay_hub.db.models import ( Group, GroupMember, IPLog, - SwarmSource, User, ) router = APIRouter(prefix="/v1/groups", tags=["groups"]) -# Swarm endpoints live at /v1/swarm/*. They were previously declared on the groups -# router with a full path, which mounted them at /v1/groups/v1/swarm/* (H7). -swarm_router = APIRouter(prefix="/v1/swarm", tags=["swarm"]) - @router.get("/mine") async def my_groups( @@ -211,113 +206,6 @@ async def list_public_groups( return {"groups": groups, "total": len(groups)} -# ── Swarm (content replication) ─────────────────────────────────────────────── - -class SwarmRegisterRequest(BaseModel): - content_hash: str # blake3 hex - endpoint: str # "<scheme>:<port>" — a port on the caller, never a host - - -# A transport and a port, and deliberately no host. The field used to be free -# text documented as "ip:port", so a caller could name *someone else's* -# address as a source; nothing dials a swarm source today, which is the only -# reason that was not already a reflection primitive. A reader learns where a -# node is from the node record, which is stamped with the address the announce -# actually came from — so a host here would be a second, weaker, answer to a -# question already settled elsewhere. -_SWARM_ENDPOINT = re.compile(r"^(webrtc|quic):([0-9]{1,5})$") - -# One account, this many public hashes. Rows are keyed (hash, account) with no -# cap, so a loop of invented hashes was unbounded storage growth on a hub -# shared with everyone else. A public library far larger than this is a real -# thing — but it is one a hub operator should be asked about, not something a -# client establishes by writing rows. -MAX_SWARM_HASHES_PER_ACCOUNT = 10_000 - - -@swarm_router.post("/register", status_code=201) -@limiter.limit("120/minute") -async def swarm_register( - body: SwarmRegisterRequest, - request: Request, - current_user: User = Depends(get_current_user), - db: AsyncSession = Depends(get_db), -): - """ - Node registers itself as a source for a PUBLIC content hash. - - Finding H7: the node registered hashes for every group it hosted, private ones - included, and this route was mounted at /v1/groups/v1/swarm/register — so the - node's calls 404'd and the leak was masked by a routing bug rather than - prevented. Nodes now filter by group visibility before calling, and the path is - correct, so the filter has to be right. - - Availability: the endpoint is a port, not an address, and the number of - hashes one account may claim is bounded. See the two constants above. - """ - from meshbay_hub.csam import check_content_hash - if check_content_hash(body.content_hash): - raise HTTPException(status_code=451, detail="Content blocked") - - m = _SWARM_ENDPOINT.match(body.endpoint or "") - if not m or not (0 < int(m.group(2)) < 65536): - raise HTTPException( - status_code=422, - detail="endpoint must be '<webrtc|quic>:<port>' — a port on the " - "registering node, not an address") - - from datetime import datetime - existing = await db.get(SwarmSource, (body.content_hash, current_user.id)) - now = datetime.now(UTC) - if existing: - existing.endpoint = body.endpoint - existing.last_seen = now - else: - held = (await db.execute( - select(func.count()).select_from(SwarmSource) - .where(SwarmSource.node_id == current_user.id))).scalar() or 0 - if held >= MAX_SWARM_HASHES_PER_ACCOUNT: - raise HTTPException( - status_code=429, - detail="This account already claims the maximum number of " - "public content hashes") - db.add(SwarmSource( - content_hash=body.content_hash, - node_id=current_user.id, - endpoint=body.endpoint, - )) - await db.commit() - return {"status": "registered", "hash": body.content_hash} - - -@swarm_router.get("/{content_hash}") -async def swarm_sources( - content_hash: str, - current_user: User = Depends(get_current_user), - db: AsyncSession = Depends(get_db), -): - """ - Return nodes that can serve a content hash. - - Authenticated (H7): an open endpoint lets anyone probe whether a given file - exists anywhere in the network and which node holds it. - """ - from datetime import datetime, timedelta - cutoff = datetime.now(UTC) - timedelta(minutes=30) - result = await db.execute( - select(SwarmSource) - .where( - SwarmSource.content_hash == content_hash, - SwarmSource.last_seen > cutoff, - ) - ) - sources = result.scalars().all() - return { - "hash": content_hash, - "sources": [{"node_id": s.node_id, "endpoint": s.endpoint} for s in sources], - } - - @router.get("/{group_id}/members") async def group_members( group_id: str, diff --git a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py index 35c688c..038f310 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py @@ -1,59 +1,88 @@ """ -MeshBay Hub — moderation endpoints. +MeshBay Hub — moderation endpoints (docs/MESHBAY_DESIGN.md §7.5). -Reporting flow: - POST /v1/reports — report a content hash (sign-in required) +Reporting: + POST /v1/reports — a member of a public group reports a file they saw there. - 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 + Who may: a person's account (never a node's token), old enough + (`reports.min_account_age_hours`), an active member of that public group, + within a daily allowance (`reports.daily_per_account`) as well as the per-address + rate limit. One report per account per hash. + + What it leads to: once `reports.review_threshold` distinct accounts have + reported a hash, it is queued for an administrator (`content_reviews`), who is + notified and blocks or dismisses it. With `reports.auto_block` on, it is blocked + at once instead — the instance's choice, off by default, because a handful of + accounts made for the purpose would then be enough to take a file down. 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. + switched off there is nothing here to report. Admin endpoints: - GET /v1/admin/blocklist — list blocked hashes - POST /v1/admin/blocklist — manually add a hash - DELETE /v1/admin/blocklist/{hash} — remove a hash + GET /v1/admin/reports — hashes waiting for a decision + POST /v1/admin/reports/{hash}/block — block it, and tell the nodes + POST /v1/admin/reports/{hash}/dismiss — close it without blocking + GET /v1/admin/blocklist — list blocked hashes + POST /v1/admin/blocklist — manually add a hash + DELETE /v1/admin/blocklist/{hash} — remove a hash Node integration: - GET /v1/blocklist/check?hash=<blake3> — check if a hash is blocked - GET /v1/blocklist — full blocklist (for node sync) + GET /v1/blocklist?after=<hash> — the list, paged, for a node's own token + WebSocket `blocklist_update` — additions and removals, pushed to nodes + hosting a public group """ import logging +from collections import Counter +from datetime import UTC, datetime, timedelta +from typing import Literal from fastapi import APIRouter, Depends, HTTPException, Query, Request -from pydantic import BaseModel +from pydantic import BaseModel, Field 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.deps import ( + require_admin, + require_moderator, + require_node_scope, + require_user_scope, + user_is_admin, +) from meshbay_hub.api.middleware import limiter from meshbay_hub.api.netutil import client_ip +from meshbay_hub.api.revocation import broadcast_blocklist_update from meshbay_hub.db.engine import get_db -from meshbay_hub.db.models import ContentBlocklist, ContentReport, User +from meshbay_hub.db.models import ( + ContentBlocklist, + ContentReport, + ContentReview, + Group, + GroupMember, + User, +) log = logging.getLogger(__name__) router = APIRouter(tags=["moderation"]) -# 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 +# One answer for every reason a report is not accepted from this account for this +# group, so the endpoint does not tell anyone which groups exist or who is in them. +_NOT_YOURS = "You can report a file only in a public group you are a member of." + + +def _is_hash(value: str) -> bool: + return len(value) == 64 and all(c in "0123456789abcdef" for c in value) # ── Models ──────────────────────────────────────────────────────────────────── class ReportRequest(BaseModel): - content_hash: str # blake3 hex (64 chars) - group_id: str | None = None - reason: str = "illegal" - detail: str | None = None + content_hash: str = Field(max_length=64) # blake3 hex (64 chars) + group_id: str = Field(max_length=36) + reason: Literal["illegal", "spam", "copyright", "other"] = "illegal" + detail: str | None = Field(default=None, max_length=256) class BlocklistAddRequest(BaseModel): @@ -61,124 +90,150 @@ class BlocklistAddRequest(BaseModel): reason: str -# ── Public endpoints ────────────────────────────────────────────────────────── +# ── Reporting ───────────────────────────────────────────────────────────────── @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), + current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): """ - Report a public content hash for moderation. + Report a file of a public group, as a member of that group. - 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. + Every bound here answers what a report costs someone else: a file taken out + of a group everyone else uses, and an administrator's time. So a report takes + a person's account (a node's token is refused), one that has existed for a + while, membership of the public group the file was seen in, and a daily + allowance per account besides the rate limit per address — an address is one + of thousands a subscriber holds. It never blocks anything by itself unless the + instance chose automatic blocking: past the threshold, an administrator + decides. """ 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): + if not _is_hash(body.content_hash): raise HTTPException(status_code=422, detail="content_hash must be 64 hex chars (blake3)") + limits = await hub_settings.report_limits(db) + now = datetime.now(UTC) + + created = current_user.created_at + if created is not None and created.tzinfo is None: + created = created.replace(tzinfo=UTC) + if created is not None and \ + now - created < timedelta(hours=limits["min_account_age_hours"]): + raise HTTPException(status_code=403, + detail="This account is too new to report content yet.") + + group = await db.get(Group, body.group_id) + member = await db.scalar(select(GroupMember.user_id).where( + GroupMember.group_id == body.group_id, + GroupMember.user_id == current_user.id)) + if group is None or group.visibility != "public" or group.status != "active" \ + or member is None: + raise HTTPException(status_code=403, detail=_NOT_YOURS) + + today = await db.scalar(select(func.count(ContentReport.id)).where( + ContentReport.reporter_id == current_user.id, + ContentReport.reported_at > now - timedelta(days=1))) or 0 + if today >= limits["daily_per_account"]: + raise HTTPException(status_code=429, + detail="You have reached today's number of reports.") + # 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)) + if already: + return {"status": "already_reported"} - 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() + 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() distinct_reporters = await db.scalar( select(func.count(func.distinct(ContentReport.reporter_id))) .where(ContentReport.content_hash == body.content_hash)) or 0 - 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( - content_hash=body.content_hash, - reason=f"auto:{body.reason}", - added_by="auto", - )) - action = "auto_blocked" + blocked_now = False + if distinct_reporters >= limits["review_threshold"] \ + and await db.get(ContentBlocklist, body.content_hash) is None: + review = await db.get(ContentReview, body.content_hash) + if review is not None and review.status == "dismissed": + pass # an administrator's decision stands; more reports do not reopen it + elif limits["auto_block"]: + db.add(ContentBlocklist(content_hash=body.content_hash, + reason=f"auto:{body.reason}", added_by="auto")) + if review is None: + db.add(ContentReview(content_hash=body.content_hash, status="blocked", + decided_at=now, decided_by="auto")) + else: + review.status, review.decided_at, review.decided_by = "blocked", now, "auto" + blocked_now = True log.warning("Content auto-blocked after %d distinct reporters: %s", distinct_reporters, body.content_hash[:16]) - + elif review is None: + db.add(ContentReview(content_hash=body.content_hash, status="pending")) + await _notify_admins(db, body.content_hash) + log.warning("Content queued for review after %d distinct reporters: %s", + distinct_reporters, body.content_hash[:16]) await db.commit() - return { - "status": action, - "content_hash": body.content_hash, - "report_count": distinct_reporters, - "threshold": AUTO_BLOCK_THRESHOLD, - } - + if blocked_now: + await broadcast_blocklist_update(db, add=[body.content_hash]) + # The same answer whatever happened next: a reporter is not told how close a + # file is to review, which is a count to aim at. + return {"status": "logged"} -@router.get("/v1/blocklist/check") -@limiter.limit("120/minute") -async def check_blocklist( - hash: str, - request: Request, - db: AsyncSession = Depends(get_db), -): - """Check if a single hash is blocked. Used by nodes before serving public content. - Unauthenticated, because a node consults it before serving public content - and does so on its own behalf. That makes the shape check worth having: - without it any string of any length became a primary-key lookup. - """ - if len(hash) != 64 or not all(c in "0123456789abcdef" for c in hash): - raise HTTPException(status_code=422, detail="hash must be 64 hex chars (blake3)") - blocked = await db.get(ContentBlocklist, hash) - return { - "blocked": blocked is not None, - "hash": hash, - "reason": blocked.reason if blocked else None, - } +async def _notify_admins(db: AsyncSession, content_hash: str) -> None: + from meshbay_hub.api.notifications import create_notification + admins = [u for u in (await db.execute(select(User).where( + User.status == "active"))).scalars().all() if user_is_admin(u)] + for admin in admins: + await create_notification( + db, admin.id, "content_review", + "Reported content is waiting for a decision", + detail=content_hash[:16], link="#/admin", aggregate=False) @router.get("/v1/blocklist") async def get_blocklist( + current_node: User = Depends(require_node_scope), db: AsyncSession = Depends(get_db), - # Bounded, like every other list. This one takes no authentication — a - # node syncs it at startup — and had no ceiling at all, so any stranger - # could ask for the table in one query, repeatedly. 10 000 is what a node - # asks for, so it is the default and also the most anyone may have. + after: str = Query(default="", max_length=64), limit: int = Query(default=10000, ge=1, le=10000), ): """ - Return the full blocklist. Nodes sync this on startup. - Returns hashes only (not reasons) to minimize data exposure. + The content blocklist, a page at a time, for a node hosting a public group. + + Hashes only, never the reasons. Ordered by hash so `after` (the last hash of + the previous page) is a stable cursor: a list longer than one page used to + be cut at 10 000 with no way to ask for the rest, and the node applying it + silently served everything past the cut. A node's own token, because this + is what a node fetches on its own behalf and nothing else asks for it. """ result = await db.execute( select(ContentBlocklist.content_hash) - .order_by(ContentBlocklist.added_at.desc()) + .where(ContentBlocklist.content_hash > after) + .order_by(ContentBlocklist.content_hash) .limit(limit) ) hashes = [row[0] for row in result.fetchall()] - return {"count": len(hashes), "hashes": hashes} + return {"hashes": hashes, + "next": hashes[-1] if len(hashes) == limit else None} # ── Admin endpoints ─────────────────────────────────────────────────────────── @@ -214,16 +269,19 @@ async def admin_add_blocklist( current_user: User = Depends(require_admin), db: AsyncSession = Depends(get_db), ): + if not _is_hash(body.content_hash): + raise HTTPException(status_code=422, detail="content_hash must be 64 hex chars (blake3)") existing = await db.get(ContentBlocklist, body.content_hash) if existing: raise HTTPException(status_code=409, detail="Hash already blocked") db.add(ContentBlocklist( content_hash=body.content_hash, - reason=body.reason, + reason=body.reason[:64], added_by=current_user.username, )) await db.commit() + await broadcast_blocklist_update(db, add=[body.content_hash]) return {"status": "blocked", "hash": body.content_hash} @@ -238,5 +296,73 @@ async def admin_remove_blocklist( raise HTTPException(status_code=404, detail="Hash not in blocklist") await db.delete(entry) await db.commit() + await broadcast_blocklist_update(db, remove=[content_hash]) return {"status": "unblocked", "hash": content_hash} + +# ── Review queue ────────────────────────────────────────────────────────────── + +@router.get("/v1/admin/reports") +async def admin_list_reports( + current_user: User = Depends(require_moderator), + db: AsyncSession = Depends(get_db), + limit: int = Query(default=100, ge=1, le=500), +): + """Hashes waiting for a decision, oldest first, with what was said about them.""" + reviews = (await db.execute( + select(ContentReview).where(ContentReview.status == "pending") + .order_by(ContentReview.opened_at).limit(limit))).scalars().all() + out = [] + for r in reviews: + reports = (await db.execute(select(ContentReport).where( + ContentReport.content_hash == r.content_hash))).scalars().all() + group_ids = sorted({x.group_id for x in reports if x.group_id}) + names = dict((await db.execute(select(Group.id, Group.name).where( + Group.id.in_(group_ids)))).all()) if group_ids else {} + out.append({ + "hash": r.content_hash, + "opened_at": r.opened_at.isoformat(), + "reporters": len({x.reporter_id for x in reports}), + "reasons": dict(Counter(x.reason for x in reports)), + "details": [x.detail for x in reports if x.detail][:10], + "groups": [{"id": g, "name": names.get(g, "")} for g in group_ids], + }) + return {"reports": out} + + +async def _decide(db: AsyncSession, content_hash: str, status: str, by: str) -> ContentReview: + review = await db.get(ContentReview, content_hash) + if review is None or review.status != "pending": + raise HTTPException(status_code=404, detail="Nothing waiting for this hash") + review.status, review.decided_at, review.decided_by = status, datetime.now(UTC), by + return review + + +@router.post("/v1/admin/reports/{content_hash}/block") +async def admin_block_reported( + content_hash: str, + current_user: User = Depends(require_admin), + db: AsyncSession = Depends(get_db), +): + await _decide(db, content_hash, "blocked", current_user.username) + reasons = Counter((await db.execute(select(ContentReport.reason).where( + ContentReport.content_hash == content_hash))).scalars().all()) + if await db.get(ContentBlocklist, content_hash) is None: + db.add(ContentBlocklist( + content_hash=content_hash, + reason=f"reported:{reasons.most_common(1)[0][0] if reasons else 'other'}", + added_by=current_user.username)) + await db.commit() + await broadcast_blocklist_update(db, add=[content_hash]) + return {"status": "blocked", "hash": content_hash} + + +@router.post("/v1/admin/reports/{content_hash}/dismiss") +async def admin_dismiss_reported( + content_hash: str, + current_user: User = Depends(require_admin), + db: AsyncSession = Depends(get_db), +): + await _decide(db, content_hash, "dismissed", current_user.username) + await db.commit() + return {"status": "dismissed", "hash": content_hash} diff --git a/packages/meshbay-hub/src/meshbay_hub/api/relay.py b/packages/meshbay-hub/src/meshbay_hub/api/relay.py deleted file mode 100644 index 7bb3f66..0000000 --- a/packages/meshbay-hub/src/meshbay_hub/api/relay.py +++ /dev/null @@ -1,166 +0,0 @@ -""" -MeshBay Hub — Mesh Relay registration protocol (5.3). - -Community-operated TURN relays register with hubs. -Nodes query the hub for available relays when UDP hole punching fails. - -Relay registration: - POST /v1/relays/register — relay announces itself (signed JWT) - GET /v1/relays — list active relays (for nodes) - -Relay authentication: relay generates an Ed25519 keypair at install time. An -admin approves the public key, and every register call carries an Ed25519 -signature over "meshbay:relay_register:<relay_id>:<endpoint>:<timestamp>" — -the same proof-of-possession shape as /v1/nodes/announce. - -Relay is responsible for E2E encrypted QUIC traffic only (it cannot -read the application-layer content, only forward UDP packets). -""" - -import base64 -import logging -import time - -from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey -from fastapi import APIRouter, Depends, HTTPException -from pydantic import BaseModel -from sqlalchemy.ext.asyncio import AsyncSession - -from meshbay_hub.api.deps import require_admin -from meshbay_hub.db.engine import get_db -from meshbay_hub.db.models import User - -log = logging.getLogger(__name__) - -# **Closed, the same way and for a similar reason as federation.** Nothing in the -# tree calls these routes — no node asks for a relay, no client offers one — and -# §11.1 measured two ISPs with no TURN relay needed. Two of the three take no -# account and answer anyone who can reach the hub, so a registry nothing uses -# was an unauthenticated surface kept for its own sake. A constant, not a -# setting: re-opening it means building the node side first, then flipping this. -RELAYS_ENABLED = False - - -def _relays_open() -> None: - """Refuse every route on this router while the registry is closed. - - On the router rather than in each handler, so a route added later is closed - before anybody remembers to write the check (C6). - """ - if not RELAYS_ENABLED: - raise HTTPException(status_code=503, - detail="The relay registry is not enabled on this hub") - - -router = APIRouter(prefix="/v1/relays", tags=["relay"], - dependencies=[Depends(_relays_open)]) - -# In-memory relay registry (production: DB table) -_relays: dict[str, dict] = {} # relay_id → {endpoint, pk, last_seen, capacity} - - -# ── Models ──────────────────────────────────────────────────────────────────── - -class RelayRegisterRequest(BaseModel): - """Relay self-registers, proving possession of its approved key.""" - relay_id: str - endpoint: str # "ip:port" (UDP) - pk_relay: str # base64 Ed25519 public key - capacity: int = 100 # max concurrent connections - timestamp: int | None = None # unix seconds - signature: str | None = None # base64 Ed25519 over the register message - - -class RelayAdminApproveRequest(BaseModel): - relay_id: str - pk_relay: str # admin approves by registering the relay's public key - - -# ── Relay endpoints ─────────────────────────────────────────────────────────── - -REGISTER_TIMESTAMP_WINDOW = 300 # seconds either side, as /v1/nodes/announce - - -@router.post("/register", status_code=201) -async def relay_register( - body: RelayRegisterRequest, - db: AsyncSession = Depends(get_db), -): - """ - Relay announces itself. Must be pre-approved by a hub admin, and must prove - it holds the private key that approval registered. - - This endpoint has no `Depends` on an account on purpose — a relay is not a - user — but it had no proof of anything either: it compared `pk_relay` - against the approved value, which is a **public** key, so anyone who could - read it could rewrite where the hub tells nodes to send relayed traffic. - The module docstring said "signs keepalive JWTs" and nothing verified a - signature; `jwt` was imported and never used. A key is not a password, and - the fix is the proof-of-possession pattern already used by - /v1/nodes/announce and /v1/nodes/auth. - """ - approved = _relays.get(body.relay_id) - if not approved or approved.get("pk") != body.pk_relay: - raise HTTPException(status_code=403, - detail="Relay not approved — ask the hub admin to " - "run POST /v1/relays/approve") - - if body.timestamp is None or not body.signature: - raise HTTPException( - status_code=400, - detail="register requires timestamp and signature (proof of possession)") - if abs(int(time.time()) - body.timestamp) > REGISTER_TIMESTAMP_WINDOW: - raise HTTPException(status_code=401, detail="Timestamp too old or too far ahead") - - message = (f"meshbay:relay_register:{body.relay_id}:" - f"{body.endpoint}:{body.timestamp}").encode() - try: - pk = Ed25519PublicKey.from_public_bytes(base64.b64decode(body.pk_relay)) - pk.verify(base64.b64decode(body.signature), message) - except Exception: - log.warning("Relay %s failed proof of possession", body.relay_id[:8]) - raise HTTPException(status_code=401, detail="Invalid relay key proof of possession") - - _relays[body.relay_id].update({ - "endpoint": body.endpoint, - "capacity": body.capacity, - "last_seen": int(time.time()), - "active": True, - }) - log.info("Relay registered: %s at %s", body.relay_id[:8], body.endpoint) - return {"status": "registered", "relay_id": body.relay_id} - - -@router.get("") -async def list_relays(): - """ - List active Mesh Relays. Called by nodes when UDP hole punching fails. - Returns only active relays (seen in the last 5 minutes). - """ - cutoff = int(time.time()) - 300 - active = [ - { - "relay_id": rid, - "endpoint": r["endpoint"], - "capacity": r["capacity"], - } - for rid, r in _relays.items() - if r.get("active") and r.get("last_seen", 0) > cutoff - ] - return {"relays": active, "count": len(active)} - - -@router.post("/approve", status_code=201) -async def admin_approve_relay( - body: RelayAdminApproveRequest, - current_user: User = Depends(require_admin), -): - """Admin: pre-approve a relay by registering its public key.""" - _relays[body.relay_id] = { - "pk": body.pk_relay, - "approved_by": current_user.username, - "approved_at": int(time.time()), - "active": False, # becomes True after first register call - } - log.info("Relay approved by %s: %s", current_user.username, body.relay_id[:8]) - return {"status": "approved", "relay_id": body.relay_id} diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py index 2c0b8db..fe9edf3 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py @@ -165,6 +165,35 @@ async def broadcast_revocation(token: str) -> int: return sent +async def broadcast_blocklist_update(db: AsyncSession, *, add: list[str] = (), + remove: list[str] = ()) -> int: + """ + Tell every connected node that hosts a public group what changed on the + content blocklist. Returns how many were told. + + Only those nodes: the list names public content, and a node hosting only + private groups has nothing to apply it to. A node that is offline now syncs + the whole list when it next connects (`GET /v1/blocklist`), so a missed push + is only late, never lost. + """ + if not add and not remove: + return 0 + public = set((await db.execute( + select(Group.id).where(Group.visibility == "public"))).scalars().all()) + payload = json.dumps({"type": "blocklist_update", + "add": list(add), "remove": list(remove)}) + sent = 0 + for node_id, ws in list(_connected_nodes.items()): + if not public.intersection(_node_groups.get(node_id, [])): + continue + try: + await ws.send_text(payload) + sent += 1 + except Exception: + _connected_nodes.pop(node_id, None) + return sent + + def _sign_revocation(target: str, target_id: str, reason: str) -> str: """Issue a signed revocation token (JWT EdDSA).""" from meshbay_hub.auth import _hub_id, _hub_sk_pem diff --git a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py index fc40204..6c9699b 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py @@ -171,7 +171,7 @@ async def webrtc_offer( Finding H4: it also ignored group status, so "suspend a group" did not stop new connections from being brokered to nodes hosting it. """ - from meshbay_hub.api.revocation import _connected_nodes, _node_groups + from meshbay_hub.api.revocation import _connected_nodes if len(body.sdp) > MAX_SDP_BYTES: raise HTTPException(status_code=413, detail="SDP too large") diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py index 7046c2f..3e996a7 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/users.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py @@ -41,7 +41,6 @@ from meshbay_hub.db.models import ( Node, Notification, RefreshToken, - SwarmSource, User, UserDevice, UserPreference, @@ -1379,14 +1378,13 @@ async def erase_account(db: AsyncSession, user: User, owned_groups: str = "refus the person it is about. Gone: credentials, email, node key, group memberships, notifications, refresh - tokens, node registrations, device keys, public-swarm sources. The username + tokens, node registrations, device keys. The username is released. Device keys go even though the desktop client keeps its private half: left behind, the key still belongs to this tombstone, so an account created later from the same installation is refused that device ("belongs to another - account"). Swarm sources are keyed by the *user* id and carry the node's - ip:port. + account"). Kept: the row itself, emptied, and the IP log that points at it. Those logs exist for one year to answer legal requests, and a log that cannot say whose @@ -1419,7 +1417,6 @@ async def erase_account(db: AsyncSession, user: User, owned_groups: str = "refus await db.execute(delete(RefreshToken).where(RefreshToken.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(SwarmSource).where(SwarmSource.node_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 |