aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/deps.py11
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/moderation.py68
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py29
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