diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/hub_client.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/hub_client.py | 43 |
1 files changed, 26 insertions, 17 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py index 346a4cd..2ce8993 100644 --- a/packages/meshbay-node/src/meshbay_node/hub_client.py +++ b/packages/meshbay-node/src/meshbay_node/hub_client.py @@ -6,7 +6,6 @@ Handles all communication from the node to a Mesh Hub: - JWT offline verification and auto-refresh - Node announcement (endpoint_hint) - User public key lookup (for GEK wrapping) - - Swarm hash registration The node authenticates via Ed25519 challenge-response (/v1/nodes/auth). No auth_key or password is ever stored on or transmitted from the node. @@ -335,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. @@ -394,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) @@ -406,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 @@ -461,26 +471,25 @@ class HubClient: log.warning("Could not deliver WebRTC answer to %s: %s", str(msg.get("peer_id"))[:8], e) - # ── Swarm registration ───────────────────────────────────────────────── + # ── Content blocklist (public groups) ───────────────────────────────── - async def register_swarm(self, content_hashes: list[str], endpoint: str) -> int: - """Register file hashes in the hub swarm table. Returns count registered.""" + 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() - - registered = 0 - for h in content_hashes: - try: - r = await self._http.post("/v1/swarm/register", json={ - "content_hash": h, - "endpoint": endpoint, - }, headers=self._session.auth_headers) - if r.status_code in (201, 200): - registered += 1 - except Exception: - pass - return registered + 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 ─────────────────────────────────── |