aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/hub_client.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/hub_client.py')
-rw-r--r--packages/meshbay-node/src/meshbay_node/hub_client.py31
1 files changed, 31 insertions, 0 deletions
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: