"""REST client in plain Python 3.8+ (standard library only). CYBERPLEX_URL=https://cyberplex.replit.app CYBERPLEX_API_KEY=ak_... python python_rest.py Optional: RUN_SECONDS (how long to keep listening; default 2) and CYBERPLEX_TOKEN_TTL_SEC (default 600). """ import json import os import threading import time import urllib.error import urllib.parse import urllib.request URL = os.environ.get("CYBERPLEX_URL", "https://cyberplex.replit.app") API_KEY = os.environ.get("CYBERPLEX_API_KEY") TTL_SEC = int(os.environ.get("CYBERPLEX_TOKEN_TTL_SEC", "600")) RUN_SECONDS = float(os.environ.get("RUN_SECONDS", "2")) if not API_KEY: raise SystemExit("set CYBERPLEX_API_KEY to your tenant API key") CHANNEL = "python-demo" class ApiError(Exception): def __init__(self, status, code, message, retry_after=None): super().__init__(f"{status} {code} {message}") self.status, self.code, self.retry_after = status, code, retry_after def call(token, method, path, body=None, timeout=40): data = json.dumps(body).encode() if body is not None else None req = urllib.request.Request(URL + path, data=data, method=method) req.add_header("Authorization", f"Bearer {token}") if data is not None: req.add_header("Content-Type", "application/json") try: with urllib.request.urlopen(req, timeout=timeout) as res: return json.load(res) except urllib.error.HTTPError as e: try: err = json.load(e).get("error", {}) except Exception: err = {} err = err if isinstance(err, dict) else {"message": str(err)} retry = e.headers.get("Retry-After") raise ApiError(e.code, err.get("code", ""), err.get("message", ""), int(retry) if retry else None) # Agent tokens are short-lived, so a long-running client must renew them. Keep one token, renew it at 80% of its life, and if the # server still answers 401 (clock skew, or the service rotated its signing secret) mint a fresh one and retry once. The API key does # not expire on its own, so it can always mint again. Do this minting on your SERVER, never in a browser (it needs the secret key). _token = {"value": None, "renew_at": 0.0} _lock = threading.Lock() # two threads (the reader and the publisher) share the token def mint(): r = call(API_KEY, "POST", "/v1/agent-token", {"agentId": "python-agent", "ttlSec": TTL_SEC}) if _token["value"]: print("(renewing token)") _token["value"] = r["token"] _token["renew_at"] = time.time() + r["expiresInSec"] * 0.8 def api(method, path, body=None, timeout=40): with _lock: if not _token["value"] or time.time() >= _token["renew_at"]: mint() tok = _token["value"] try: return call(tok, method, path, body, timeout) except ApiError as e: if e.status != 401: raise print("(token rejected: minting a new one)") with _lock: if _token["value"] == tok: # nobody else already replaced it mint() tok = _token["value"] return call(tok, method, path, body, timeout) stop = threading.Event() def subscribe(channel, on_message): """Long-poll loop: ask for the head cursor once, then resume from the last cursor seen.""" cursor = api("GET", f"/v1/channels/{channel}/messages?limit=1")["cursor"] while not stop.is_set(): try: q = urllib.parse.urlencode({"cursor": cursor, "wait": 20, "limit": 100}) page = api("GET", f"/v1/channels/{channel}/messages?{q}") for m in page["publications"]: on_message(m["data"]) cursor = m["cursor"] # advance per message if not page["publications"]: cursor = page["cursor"] except ApiError as e: time.sleep(e.retry_after or 2) # 429: wait as told; the cursor is kept so nothing is lost except OSError: time.sleep(2) # a network blip or a restart reader = threading.Thread(target=subscribe, args=(CHANNEL, lambda m: print("received:", json.dumps(m))), daemon=True) reader.start() time.sleep(0.5) # Publish a few messages. for n in (1, 2, 3): cursor = api("POST", f"/v1/channels/{CHANNEL}/messages", {"data": {"n": n}})["cursor"] print(f"published n={n} -> cursor {cursor}") time.sleep(RUN_SECONDS) stop.set() # daemon thread: the process exits even though a long poll may still be in flight