summaryrefslogtreecommitdiffstats
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.py54
1 files changed, 54 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 d91b945..74851c1 100644
--- a/packages/meshbay-node/src/meshbay_node/hub_client.py
+++ b/packages/meshbay-node/src/meshbay_node/hub_client.py
@@ -19,6 +19,7 @@ import logging
import time
from dataclasses import dataclass, field
from pathlib import Path
+from typing import Any, Callable
import httpx
import jwt
@@ -243,6 +244,59 @@ class HubClient:
r.raise_for_status()
return r.json()
+ # ── Persistent WebSocket (signaling + revocations) ──────────────────────
+
+ async def maintain_ws(
+ self,
+ on_incoming: Any = None,
+ on_revocation: Any = None,
+ ) -> None:
+ """
+ Maintain a persistent WebSocket connection to the hub.
+ Receives NAT punch requests and revocation tokens.
+ Runs until cancelled.
+ """
+ import websockets
+
+ if self._session is None:
+ raise RuntimeError("Not logged in")
+
+ hub_url = self._session.hub_url.replace("https://", "wss://").replace("http://", "ws://")
+ ws_url = f"{hub_url}/v1/nodes/ws"
+
+ while True:
+ try:
+ async with websockets.connect(ws_url) as ws:
+ await ws.send(json.dumps({
+ "type": "auth",
+ "token": self._session.access_token,
+ }))
+ auth_resp = json.loads(await ws.recv())
+ if auth_resp.get("type") != "auth_ok":
+ log.error("WS auth failed: %s", auth_resp)
+ return
+
+ log.info("Hub WS connected")
+
+ async for raw in ws:
+ msg = json.loads(raw)
+ mtype = msg.get("type")
+
+ if mtype == "client_incoming" and on_incoming:
+ await on_incoming(msg["peer_ip"], msg["peer_port"])
+ await ws.send(json.dumps({"type": "punch_ready"}))
+
+ elif mtype == "revocation" and on_revocation:
+ on_revocation(msg.get("token", ""))
+
+ elif mtype == "pong":
+ pass
+
+ except Exception as e:
+ log.warning("Hub WS disconnected: %s — reconnecting in 5s", e)
+ import asyncio
+ await asyncio.sleep(5)
+
# ── Convenience: full startup sequence ───────────────────────────────────
async def startup(self, endpoint_hint: str | None = None) -> HubSession: