aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub
diff options
context:
space:
mode:
Diffstat (limited to 'packages/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
-rw-r--r--packages/meshbay-hub/tests/test_moderation.py103
-rw-r--r--packages/meshbay-hub/tests/test_unauthenticated_surface.py2
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",