aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/admin.py25
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/deps.py11
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/groups.py112
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/moderation.py318
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/relay.py166
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py29
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/signaling.py2
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/users.py7
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