aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src')
-rw-r--r--packages/meshbay-node/src/meshbay_node/blocklist.py83
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py64
-rw-r--r--packages/meshbay-node/src/meshbay_node/hub_client.py31
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py19
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py12
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/wire.py16
9 files changed, 223 insertions, 8 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/blocklist.py b/packages/meshbay-node/src/meshbay_node/blocklist.py
new file mode 100644
index 0000000..29d8666
--- /dev/null
+++ b/packages/meshbay-node/src/meshbay_node/blocklist.py
@@ -0,0 +1,83 @@
+"""
+The hub's content blocklist, as this node applies it (docs/MESHBAY_DESIGN.md §7.5).
+
+Hashes of public content that the hub's moderation has blocked. A node hosting a
+public group stops serving those files there: they leave the index members are
+sent, and a request for one is refused. Nothing is deleted from the operator's
+disk — the node stops serving, and what the operator keeps is theirs to decide.
+
+Private groups are untouched. The hub never learns what a private group holds,
+so there is nothing it could have blocked in one, and nothing here reaches them.
+
+Kept on disk as well as in memory, so a node that restarts while the hub is
+unreachable does not serve again, meanwhile, what it had already stopped serving.
+The hub's list is the authority: a full sync replaces this one.
+
+It is an exact match on the content id. A file changed by one byte is another
+id, and files over the partial-hash threshold are identified by a sample of their
+bytes (§6.3) — a moderation tool, not a guarantee.
+"""
+
+import json
+import logging
+import os
+import re
+from pathlib import Path
+
+log = logging.getLogger(__name__)
+
+_HASH = re.compile(r"^[0-9a-f]{64}$")
+
+
+class ContentBlocklist:
+ def __init__(self, path: Path | None = None):
+ self._path = path
+ self._hashes: set[str] = set()
+ if path is not None:
+ try:
+ data = json.loads(path.read_text(encoding="utf-8"))
+ self._hashes = {h for h in data.get("hashes", []) if _HASH.match(h)}
+ except FileNotFoundError:
+ pass
+ except (OSError, ValueError) as e:
+ # A damaged file is not a reason to refuse to start; the next
+ # sync with the hub rewrites it.
+ log.warning("Content blocklist %s unreadable: %s", path, e)
+
+ def __contains__(self, content_hash: object) -> bool:
+ return content_hash in self._hashes
+
+ def __len__(self) -> int:
+ return len(self._hashes)
+
+ def __iter__(self):
+ return iter(self._hashes)
+
+ def replace(self, hashes) -> bool:
+ """The hub's full list. True when it differs from what was applied."""
+ new = {h for h in hashes if isinstance(h, str) and _HASH.match(h)}
+ if new == self._hashes:
+ return False
+ self._hashes = new
+ self._save()
+ return True
+
+ def apply(self, add=(), remove=()) -> bool:
+ """One pushed change. True when it changed anything."""
+ before = set(self._hashes)
+ self._hashes |= {h for h in add if isinstance(h, str) and _HASH.match(h)}
+ self._hashes -= {h for h in remove if isinstance(h, str)}
+ if self._hashes == before:
+ return False
+ self._save()
+ return True
+
+ def _save(self) -> None:
+ if self._path is None:
+ return
+ tmp = self._path.with_suffix(".tmp")
+ try:
+ tmp.write_text(json.dumps({"hashes": sorted(self._hashes)}), encoding="utf-8")
+ os.replace(tmp, self._path)
+ except OSError as e:
+ log.warning("Content blocklist %s not saved: %s", self._path, e)
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index a29aa9f..45a6b9a 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -44,6 +44,7 @@ from meshbay_common.protocol import MNP
from meshbay_node import uploads as uploads_mod
from meshbay_node.audit import RETENTION_DAYS as AUDIT_RETENTION_DAYS
from meshbay_node.audit import AuditStore
+from meshbay_node.blocklist import ContentBlocklist
from meshbay_node.bundle_store import BundleStore
from meshbay_node.chat.store import ChatStore
from meshbay_node.cli.dispatch import run, start
@@ -114,6 +115,10 @@ class NodeDaemon(EnrichmentMixin):
# Persisted so a restart does not silently un-revoke everyone (H4)
self._denylist = (
Denylist(path=config.data_dir / "denylist.json") if Denylist else None)
+ # The hub's content blocklist, applied in public groups only (§7.5).
+ # Persisted for the same reason: a restart while the hub is unreachable
+ # must not serve again what had stopped being served.
+ self._blocklist = ContentBlocklist(config.data_dir / "blocklist.json")
self._chat_stores: dict[str, ChatStore] = {}
# One instance, shared by every group's DirectoryIndexer — see
# indexer/cache.py's docstring for why this stopped being per-group.
@@ -525,6 +530,7 @@ class NodeDaemon(EnrichmentMixin):
self._webrtc._ctx["pk_x25519_b64"] = keys.pk_x25519_b64
self._webrtc._ctx["roster"] = self._roster
+ self._webrtc._ctx["blocklist"] = self._blocklist
# The MNP adapter calls the same operations as the loopback API
# (meshbay_node.ops), and those take the daemon's state. Handing
# the transport a second set of lookups is how two paths to one
@@ -611,11 +617,16 @@ class NodeDaemon(EnrichmentMixin):
except Exception as e:
log.warning("Invalid revocation token: %s", e)
+ async def on_connected():
+ await self._sync_blocklist(hub)
+
ws_task = asyncio.create_task(hub.maintain_ws(
on_incoming=on_incoming,
on_revocation=on_revocation,
on_webrtc_offer=on_webrtc_offer,
group_ids=lambda: list((self._state.get("groups_ctx") or {}).keys()),
+ on_connected=on_connected,
+ on_blocklist=self._on_blocklist_update,
))
self._tasks.append(ws_task)
log.info("Hub WS task started")
@@ -1448,9 +1459,10 @@ class NodeDaemon(EnrichmentMixin):
# node refuses every handshake while the GEK is None (NS8) — so this is
# "nobody is listening", not a case to send in clear for.
if peers and idx.gek:
- msg = (index_delta_message(idx, delta, indexer.roots)
+ hidden = self._blocklist if self._is_public(group_id) else ()
+ msg = (index_delta_message(idx, delta, indexer.roots, hidden)
if delta is not None
- else index_sync_message(idx, indexer.roots))
+ else index_sync_message(idx, indexer.roots, hidden))
pushed = 0
for session in peers:
try:
@@ -1462,6 +1474,54 @@ class NodeDaemon(EnrichmentMixin):
log.info("Index %s pushed to %d WebRTC peers",
"delta" if delta is not None else "sync", pushed)
+ # ── Content blocklist (public groups, §7.5) ──────────────────────────────
+
+ def _is_public(self, group_id: str) -> bool:
+ return any(g.id == group_id and g.visibility == "public"
+ for g in self._config.groups)
+
+ def _hosts_public_group(self) -> bool:
+ return any(g.visibility == "public" for g in self._config.groups)
+
+ async def _sync_blocklist(self, hub) -> None:
+ """The hub's whole list, on every connection. Only a node hosting a
+ public group asks: nothing else here could have been blocked."""
+ if not self._hosts_public_group():
+ return
+ try:
+ hashes = await hub.fetch_blocklist()
+ except Exception as e:
+ # Keep applying the list on disk; the next connection tries again.
+ log.warning("Content blocklist sync failed: %s", e)
+ return
+ if self._blocklist.replace(hashes):
+ log.info("Content blocklist: %d hashes", len(self._blocklist))
+ self._push_public_indexes()
+
+ def _on_blocklist_update(self, add: list, remove: list) -> None:
+ if not self._hosts_public_group():
+ return
+ if self._blocklist.apply(add, remove):
+ log.info("Content blocklist updated: +%d -%d", len(add), len(remove))
+ self._push_public_indexes()
+
+ def _push_public_indexes(self) -> None:
+ """Resend each public group's whole index, so what the list now hides
+ leaves every connected member's view, and what it released comes back."""
+ if not self._webrtc:
+ return
+ for indexer in self._indexers:
+ idx = indexer.index
+ if not self._is_public(idx.group_id) or not idx.gek:
+ continue
+ msg = index_sync_message(idx, indexer.roots, self._blocklist)
+ for session in list(self._webrtc._sessions.values()):
+ if session._group_id == idx.group_id:
+ try:
+ session._send(msg)
+ except Exception:
+ pass
+
def _drop_group_sessions(self, group_id: str) -> None:
"""Close live sessions for a revoked group (H4)."""
if not self._webrtc or not group_id:
diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py
index 5f47090..2ce8993 100644
--- a/packages/meshbay-node/src/meshbay_node/hub_client.py
+++ b/packages/meshbay-node/src/meshbay_node/hub_client.py
@@ -334,6 +334,8 @@ class HubClient:
on_revocation: Any = None,
on_webrtc_offer: Any = None,
group_ids: list[str] | None = None, # static list or callable returning one
+ on_connected: Any = None,
+ on_blocklist: Any = None,
) -> None:
"""
Maintain a persistent WebSocket connection to the hub.
@@ -393,6 +395,12 @@ class HubClient:
self._ws = ws
log.info("Hub WS connected")
+ if on_connected:
+ # Off the read loop, for the reason offers are: a slow
+ # fetch must not stop this socket being read.
+ task = asyncio.create_task(on_connected())
+ pending.add(task)
+ task.add_done_callback(pending.discard)
async for raw in ws:
msg = json.loads(raw)
@@ -405,6 +413,9 @@ class HubClient:
elif mtype == "revocation" and on_revocation:
on_revocation(msg.get("token", ""))
+ elif mtype == "blocklist_update" and on_blocklist:
+ on_blocklist(msg.get("add") or [], msg.get("remove") or [])
+
elif mtype == "webrtc_offer" and on_webrtc_offer:
# Answered off the read loop on purpose. Awaiting the
# handler here meant one slow negotiation stopped the
@@ -460,6 +471,26 @@ class HubClient:
log.warning("Could not deliver WebRTC answer to %s: %s",
str(msg.get("peer_id"))[:8], e)
+ # ── Content blocklist (public groups) ─────────────────────────────────
+
+ async def fetch_blocklist(self, max_pages: int = 100) -> set[str]:
+ """The hub's content blocklist, every page of it (docs/MESHBAY_DESIGN.md §7.5)."""
+ if self._session is None:
+ raise RuntimeError("Not logged in")
+ await self.ensure_fresh_token()
+ hashes: set[str] = set()
+ after = ""
+ for _ in range(max_pages):
+ r = await self._http.get("/v1/blocklist", params={"after": after},
+ headers=self._session.auth_headers)
+ r.raise_for_status()
+ page = r.json()
+ hashes.update(page.get("hashes") or [])
+ after = page.get("next") or ""
+ if not after:
+ return hashes
+ raise RuntimeError("Content blocklist longer than expected; not applied")
+
# ── Convenience: full startup sequence ───────────────────────────────────
async def startup(self, endpoint_hint: str | None = None) -> HubSession:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
index f2b6fa5..857db33 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
@@ -64,6 +64,8 @@ class MusicMixin:
"""
ctx = self._group_ctx()
file_id = msg.get("file_id", "")
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
self._send({"type": "error", "detail": "File not found"})
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py
index a425c65..4337e24 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py
@@ -228,6 +228,8 @@ class StreamingMixin:
async def _stream_video_inner(self, msg: dict) -> None:
ctx = self._group_ctx()
file_id = msg.get("file_id", "")
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
self._send({"type": "error", "detail": "File not found"})
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
index bd2ab52..70781eb 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
@@ -39,6 +39,8 @@ class SubtitlesMixin:
"""
ctx = self._group_ctx()
file_id = msg.get("file_id", "")
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
self._send({"type": "error", "detail": "File not found"})
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
index c3fb623..882ddbc 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
@@ -246,6 +246,25 @@ class SessionCore:
return self._ctx["groups"].get(self._group_id) or {}
return self._ctx
+ def _hidden_ids(self):
+ """The content blocklist, in a public group; nothing anywhere else.
+
+ A private group's content never reaches the hub, so nothing in it can have
+ been blocked there (docs/MESHBAY_DESIGN.md §7.5, `meshbay_node.blocklist`).
+ """
+ if self._group_ctx().get("visibility") != "public":
+ return ()
+ return self._ctx.get("blocklist") or ()
+
+ def _refuse_blocked(self, file_id) -> bool:
+ """Refuse a file the blocklist names, in a public group. True if refused."""
+ if not isinstance(file_id, str) or file_id not in self._hidden_ids():
+ return False
+ self._send({"type": "error", "detail": "This file is not available here.",
+ "code": "content_blocked", "file_id": file_id})
+ self._audit("content_blocked", file_id[:16])
+ return True
+
def _register_peer(self) -> None:
"""Add this connection to its group's peer set.
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py
index e76e23e..f629563 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py
@@ -179,7 +179,7 @@ class FilesMixin:
def _do_index_sync(self) -> None:
ctx = self._group_ctx()
- self._send(index_sync_message(ctx["index"], ctx.get("roots")))
+ self._send(index_sync_message(ctx["index"], ctx.get("roots"), self._hidden_ids()))
async def _try_serve_thumbnail(
self, thumb_hash: str, chunk_index: int, gek: bytes | None,
@@ -236,8 +236,18 @@ class FilesMixin:
self._note_unleased(tr)
file_id = msg["file_id"]
chunk_index = msg["chunk_index"]
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
+ # A blocked file's own thumbnail is a preview of it.
+ hidden = self._hidden_ids()
+ if hidden and file_id in {
+ e.thumb_hash for e in map(ctx["index"].get_entry, hidden)
+ if e is not None and e.thumb_hash}:
+ self._send({"type": "error", "detail": "This file is not available here.",
+ "code": "content_blocked", "file_id": file_id})
+ return
thumb = await self._try_serve_thumbnail(file_id, chunk_index, ctx.get("gek"))
if thumb is not None:
log.debug("file_req file_id=%s chunk=%s: served as thumbnail",
diff --git a/packages/meshbay-node/src/meshbay_node/transport/wire.py b/packages/meshbay-node/src/meshbay_node/transport/wire.py
index 1f9c925..258edb5 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/wire.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/wire.py
@@ -65,7 +65,7 @@ def list_dirs(roots: RootSet | None) -> list[str]:
return sorted(out)[:MAX_DIRS]
-def index_sync_message(index, roots: RootSet | None) -> dict:
+def index_sync_message(index, roots: RootSet | None, hidden=()) -> dict:
"""
The full `index_sync` message for one group.
@@ -73,10 +73,13 @@ def index_sync_message(index, roots: RootSet | None) -> dict:
without them a folder someone just created, or one they emptied, does not exist as
far as a client is concerned, and a member cannot tell "the drive is unplugged"
from "it is all still there".
+
+ `hidden` is the content blocklist in a public group (`meshbay_node.blocklist`):
+ those entries are not sent at all.
"""
payload = {
"version": index.version,
- "entries": [index_entry_wire(e) for e in index.entries],
+ "entries": [index_entry_wire(e) for e in index.entries if e.id not in hidden],
"dirs": list_dirs(roots),
"roots": roots.describe() if roots else [],
}
@@ -88,7 +91,7 @@ def index_sync_message(index, roots: RootSet | None) -> dict:
}
-def index_delta_message(index, delta, roots=None) -> dict:
+def index_delta_message(index, delta, roots=None, hidden=()) -> dict:
"""
One `index_delta` — what changed since the last thing this node broadcast.
@@ -103,13 +106,16 @@ def index_delta_message(index, delta, roots=None) -> dict:
page: the delta that told them something had changed was the one message
that could not say what. It is a handful of dicts, bounded by the number of
directories a group has, and it is sealed with the rest.
+
+ `hidden`, as for `index_sync_message`: a blocked entry is never added or
+ updated; its deletion still goes out.
"""
payload = {
"base_version": delta.base_version,
"version": delta.version,
- "additions": [index_entry_wire(e) for e in delta.additions],
+ "additions": [index_entry_wire(e) for e in delta.additions if e.id not in hidden],
"deletions": list(delta.deletions),
- "updates": [index_entry_wire(e) for e in delta.updates],
+ "updates": [index_entry_wire(e) for e in delta.updates if e.id not in hidden],
}
if roots is not None:
payload["roots"] = roots.describe()