aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/daemon.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/daemon.py')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py90
1 files changed, 61 insertions, 29 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index 220a908..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")
@@ -667,12 +678,6 @@ class NodeDaemon(EnrichmentMixin):
await indexer.initial_scan()
log.info("Background scan complete for %s: %d files",
name, indexer.index.count)
- # Swarm registration for public groups (after files are known).
- if gctx.get("visibility") == "public":
- endpoint = f"webrtc:{self._config.node.quic_port}"
- hashes = [e.id for e in gctx["index"].entries]
- if hashes:
- await self._register_swarm(hashes, endpoint)
# initial_scan() itself never calls on_change (it predates
# the concept — every existing caller only cared about the
# scan finishing, not about notifying anyone) — but Videos
@@ -1454,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:
@@ -1468,20 +1474,53 @@ class NodeDaemon(EnrichmentMixin):
log.info("Index %s pushed to %d WebRTC peers",
"delta" if delta is not None else "sync", pushed)
- # 11.9 — Register file hashes with hub swarm table (public groups only, H7)
- group_cfg = next(
- (g for g in self._config.groups if g.id == group_id), None)
- if (self._hub and self._state.get("endpoint_hint")
- and group_cfg and group_cfg.visibility == "public"):
- # Only the newly added hashes once there is a delta to know them
- # from — registering the whole library again on every change is
- # the same O(changes x library size) cost the delta above exists
- # to avoid.
- hashes = ([e.id for e in delta.additions] if delta is not None
- else [e.id for e in idx.entries])
- if hashes:
- endpoint = f"webrtc:{self._config.node.quic_port}"
- spawn(self._register_swarm(hashes, endpoint))
+ # ── 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)."""
@@ -1492,13 +1531,6 @@ class NodeDaemon(EnrichmentMixin):
spawn(session.close())
log.info("Dropped session for revoked group %s", group_id[:8])
- async def _register_swarm(self, hashes: list[str], endpoint: str) -> None:
- try:
- n = await self._hub.register_swarm(hashes, endpoint)
- log.info("Swarm: registered %d/%d hashes", n, len(hashes))
- except Exception as e:
- log.warning("Swarm registration failed: %s", e)
-
async def _shutdown(self) -> None:
log.info("Shutting down...")
self._state["status"] = "stopping"