aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/daemon.py
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-28 21:49:14 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-28 21:49:14 +0200
commit07480eb3f8ad0bb4369ac8c41df7c4140b108d0e (patch)
tree2f692c0dccb184890fd82339590a4e9ce168c6e4 /packages/meshbay-node/src/meshbay_node/daemon.py
parent91505face56f7ee6817408e52bad7902add75f09 (diff)
downloadmeshbay-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-node/src/meshbay_node/daemon.py')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py64
1 files changed, 62 insertions, 2 deletions
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: