diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-09-28 21:49:14 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-09-28 21:49:14 +0200 |
| commit | 07480eb3f8ad0bb4369ac8c41df7c4140b108d0e (patch) | |
| tree | 2f692c0dccb184890fd82339590a4e9ce168c6e4 /packages/meshbay-hub | |
| parent | 91505face56f7ee6817408e52bad7902add75f09 (diff) | |
| download | meshbay-07480eb3f8ad0bb4369ac8c41df7c4140b108d0e.tar.gz | |
feat: nodes apply the content blocklist in their public groups
A node hosting a public group syncs the hub's blocklist on every
connection (paged, node token only) and applies pushed changes. A
blocked file leaves the index and is refused (content_blocked); private
groups are untouched. The unused per-hash check route is gone.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-hub')
5 files changed, 166 insertions, 47 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 diff --git a/packages/meshbay-hub/tests/test_moderation.py b/packages/meshbay-hub/tests/test_moderation.py index 3109e29..b328d67 100644 --- a/packages/meshbay-hub/tests/test_moderation.py +++ b/packages/meshbay-hub/tests/test_moderation.py @@ -22,6 +22,20 @@ async def _register_and_login(client, username: str) -> dict: return {"Authorization": f"Bearer {r.json()['access_token']}"} +async def _node_headers(client, username: str) -> dict: + """A node daemon's token for a fresh account — what a node syncs with.""" + from meshbay_hub.auth import issue_access_token + user = await _register_and_login(client, username) + me = (await client.get("/v1/users/me", headers=user)).json() + tok = issue_access_token(me["user_id"], ttl=3600, groups=[], scope="node") + return {"Authorization": f"Bearer {tok}"} + + +async def _blocked(client, h: str) -> bool: + node = await _node_headers(client, f"node_{h[:6]}_{len(h)}") + return h in (await client.get("/v1/blocklist", headers=node)).json()["hashes"] + + @pytest.fixture async def reporter(client): return await _register_and_login(client, "reporter_one") @@ -69,8 +83,7 @@ async def test_same_reporter_cannot_walk_the_threshold(client, reporter): assert r.json()["report_count"] == 1 assert r.json()["status"] == "already_reported" - check = await client.get(f"/v1/blocklist/check?hash={h}") - assert check.json()["blocked"] is False + assert not await _blocked(client, h) @pytest.mark.asyncio @@ -84,8 +97,7 @@ async def test_auto_block_on_distinct_reporters(client): assert r.json()["status"] == "auto_blocked" assert r.json()["report_count"] == 3 - check = await client.get(f"/v1/blocklist/check?hash={h}") - assert check.json()["blocked"] is True + assert await _blocked(client, h) @pytest.mark.asyncio @@ -118,14 +130,13 @@ async def test_admin_add_remove_blocklist(client, admin_headers): headers=admin_headers) assert r.status_code == 201 - r = await client.get(f"/v1/blocklist/check?hash={hash4}") - assert r.json()["blocked"] is True + assert await _blocked(client, hash4) r = await client.delete(f"/v1/admin/blocklist/{hash4}", headers=admin_headers) assert r.status_code == 200 - r = await client.get(f"/v1/blocklist/check?hash={hash4}") - assert r.json()["blocked"] is False + node = await _node_headers(client, "node_after_unblock") + assert hash4 not in (await client.get("/v1/blocklist", headers=node)).json()["hashes"] @pytest.mark.asyncio @@ -134,6 +145,80 @@ async def test_full_blocklist(client, admin_headers): await client.post("/v1/admin/blocklist", json={"content_hash": hash5, "reason": "test"}, headers=admin_headers) - r = await client.get("/v1/blocklist") + node = await _node_headers(client, "node_full") + r = await client.get("/v1/blocklist", headers=node) assert r.status_code == 200 assert hash5 in r.json()["hashes"] + + +@pytest.mark.asyncio +async def test_only_a_node_token_reads_the_blocklist(client, reporter): + assert (await client.get("/v1/blocklist")).status_code in (401, 422) + assert (await client.get("/v1/blocklist", headers=reporter)).status_code == 403 + + +@pytest.mark.asyncio +async def test_the_blocklist_pages_past_its_limit(client, admin_headers): + hashes = sorted(f"{i:064x}" for i in range(5)) + for h in hashes: + await client.post("/v1/admin/blocklist", + json={"content_hash": h, "reason": "test"}, + headers=admin_headers) + node = await _node_headers(client, "node_pager") + seen, after = [], "" + while True: + page = (await client.get(f"/v1/blocklist?limit=2&after={after}", + headers=node)).json() + seen += page["hashes"] + if not page["next"]: + break + after = page["next"] + assert [h for h in seen if h in hashes] == hashes + + +@pytest.mark.asyncio +async def test_an_admin_cannot_block_something_that_is_not_a_hash(client, admin_headers): + r = await client.post("/v1/admin/blocklist", + json={"content_hash": "x" * 5000, "reason": "test"}, + headers=admin_headers) + assert r.status_code == 422 + + +@pytest.mark.asyncio +async def test_a_change_is_pushed_to_nodes_hosting_a_public_group_only(db_session): + """A node hosting only private groups has nothing to apply the list to.""" + import json + + import meshbay_hub.api.revocation as rev + from meshbay_hub.db.models import Group, User + + db_session.add(User(id="owner-bl", username="owner_bl", email="x", hub_id="h", + pw_hash=b"x", pw_salt=b"x")) + db_session.add(Group(id="pub-bl", name="pub", admin_id="owner-bl", + visibility="public")) + db_session.add(Group(id="priv-bl", name="priv", admin_id="owner-bl", + visibility="private")) + await db_session.commit() + + class _WS: + def __init__(self): + self.sent = [] + + async def send_text(self, text): + self.sent.append(json.loads(text)) + + public_node, private_node = _WS(), _WS() + saved = dict(rev._connected_nodes), dict(rev._node_groups) + rev._connected_nodes.update({"n-pub": public_node, "n-priv": private_node}) + rev._node_groups.update({"n-pub": ["pub-bl"], "n-priv": ["priv-bl"]}) + try: + sent = await rev.broadcast_blocklist_update(db_session, add=["a" * 64]) + finally: + rev._connected_nodes.clear() + rev._connected_nodes.update(saved[0]) + rev._node_groups.clear() + rev._node_groups.update(saved[1]) + assert sent == 1 + assert public_node.sent == [{"type": "blocklist_update", "add": ["a" * 64], + "remove": []}] + assert private_node.sent == [] diff --git a/packages/meshbay-hub/tests/test_unauthenticated_surface.py b/packages/meshbay-hub/tests/test_unauthenticated_surface.py index f2da42b..78615ca 100644 --- a/packages/meshbay-hub/tests/test_unauthenticated_surface.py +++ b/packages/meshbay-hub/tests/test_unauthenticated_surface.py @@ -42,8 +42,6 @@ PUBLIC = { ("POST", "/v1/nodes/auth"): "node sign-in — Ed25519 signature over a fresh timestamp", ("WS", "/v1/nodes/ws"): "a node-scoped JWT in the first message, within a timeout", ("GET", "/v1/groups"): "the public directory — empty when public groups are off", - ("GET", "/v1/blocklist"): "nodes sync it on their own behalf — hashes only", - ("GET", "/v1/blocklist/check"): "nodes consult it on their own behalf", ("GET", "/v1/relays"): "closed: 503 while relay.RELAYS_ENABLED is False", ("POST", "/v1/relays/register"): "closed; when open, approved key + signature", ("GET", "/mhp/info"): "closed: 503 while federation.FEDERATION_ENABLED is False", |