diff options
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/deps.py | 11 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/moderation.py | 68 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/revocation.py | 29 |
3 files changed, 72 insertions, 36 deletions
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/moderation.py b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py index 35c688c..c4e6af5 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py @@ -18,8 +18,9 @@ Admin endpoints: 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 @@ -30,9 +31,10 @@ 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 get_current_user, require_admin, require_node_scope 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 @@ -47,6 +49,10 @@ router = APIRouter(tags=["moderation"]) AUTO_BLOCK_THRESHOLD = 3 +def _is_hash(value: str) -> bool: + return len(value) == 64 and all(c in "0123456789abcdef" for c in value) + + # ── Models ──────────────────────────────────────────────────────────────────── class ReportRequest(BaseModel): @@ -89,7 +95,7 @@ async def report_content( 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)") # One vote per account per hash — a single reporter must not be able to walk @@ -115,6 +121,7 @@ async def report_content( .where(ContentReport.content_hash == body.content_hash)) or 0 action = "already_reported" if already else "logged" + auto_blocked = False if distinct_reporters >= AUTO_BLOCK_THRESHOLD: existing = await db.get(ContentBlocklist, body.content_hash) if not existing: @@ -124,10 +131,13 @@ async def report_content( added_by="auto", )) action = "auto_blocked" + auto_blocked = True log.warning("Content auto-blocked after %d distinct reporters: %s", distinct_reporters, body.content_hash[:16]) await db.commit() + if auto_blocked: + await broadcast_blocklist_update(db, add=[body.content_hash]) return { "status": action, "content_hash": body.content_hash, @@ -136,49 +146,31 @@ async def report_content( } -@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, - } - - @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 +206,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 +233,6 @@ 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} 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 |