From c1ed7dff65bacbe8d08faeee64ca0662f202f309 Mon Sep 17 00:00:00 2001 From: Raman Gupta <7243222+raman325@users.noreply.github.com> Date: Mon, 27 Jul 2026 16:36:55 -0400 Subject: [PATCH] python-client: fix imports and extract a reusable scrypted_client library (#2094) * python-client: link transitive server modules required by plugin_remote plugin_remote now imports cluster_labels, cluster_setup, plugin_console, plugin_pip, and plugin_volume (plus a lazy plugin_repl import), but packages/python-client only symlinks plugin_remote, rpc, rpc_reader, and scrypted_python. As a result the client cannot be imported at all: ModuleNotFoundError: No module named 'cluster_labels' Add symlinks for the missing modules, matching the existing pattern of sharing one implementation with server/python. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01WaSWgBV1XsZxWiYM16sWME * python-client: extract reusable scrypted_client library from test.py The only way to connect to Scrypted from Python has been to copy the bootstrap out of test.py, which is not importable (module-level event loop) and had drifted from the current rpc_reader/PluginRemote APIs (writeJSON vs writeSerialized, missing ClusterSetup argument). Move the transport and connection handshake into an importable scrypted_client module: - EioRpcTransport gains a close() method, an optional injectable aiohttp session for the engine.io connection, and queues its send loop on the running loop instead of via run_coroutine_threadsafe. - connect_scrypted_client() performs login, the engine.io connect, and the getRemote handshake, with a connect timeout and consistent ScryptedConnectionError on failure. It also attaches the constructed SystemManager to remote.systemManager, the same wiring loadZip does for plugins, so PluginRemote.notify can dispatch events to systemManager.listen() callbacks. - test.py becomes a small demo of the library and exits cleanly without os._exit(); the server URL is configurable via SCRYPTED_BASE_URL. Verified live against a Scrypted server (device enumeration and OnOff state reads). Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01WaSWgBV1XsZxWiYM16sWME * python-client: bound all connect phases and tear down tasks on close Address review feedback: - close() now cancels and awaits both the send loop and the peer read loop (previously the send task was cancelled but never awaited and the read task was untracked), so closing the event loop after close() no longer risks 'Task was destroyed but it is pending' warnings. - connect_scrypted_client() stores the read-loop task on the transport and propagates read-loop failures into the pending handshake future, so a link that dies mid-handshake fails immediately with the underlying error instead of waiting out the timeout. - The timeout parameter now bounds every phase: the login POST (aiohttp ClientTimeout), the engine.io connect (asyncio.wait_for), and the wait for initial system state. - Document session ownership: login_session is borrowed and never closed; an http_session given to EioRpcTransport is owned by the transport and closed by close(). Verified live against a Scrypted server: clean run with empty stderr, plus connection-refused and bad-credential paths both raising ScryptedConnectionError. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01WaSWgBV1XsZxWiYM16sWME * python-client: lazily import plugin host modules in plugin_remote A client consuming plugin_remote (for SystemManager, DeviceManager, MediaManager) only needs the types and the engine.io/rpc bits. The plugin host modules (cluster_labels, plugin_console, plugin_pip, plugin_volume) are only used inside loadZipWrapped, so import them there -- importing plugin_remote no longer requires them, and the client directory drops those symlinks (plugin_repl was already a lazy import). cluster_setup stays: the client bootstrap constructs a ClusterSetup. Co-Authored-By: Claude Fable 5 --------- Co-authored-by: Claude Fable 5 --- packages/python-client/README.md | 48 +++++ packages/python-client/cluster_setup.py | 1 + packages/python-client/scrypted_client.py | 249 ++++++++++++++++++++++ packages/python-client/test.py | 146 ++----------- server/python/plugin_remote.py | 11 +- 5 files changed, 318 insertions(+), 137 deletions(-) create mode 100644 packages/python-client/README.md create mode 120000 packages/python-client/cluster_setup.py create mode 100644 packages/python-client/scrypted_client.py diff --git a/packages/python-client/README.md b/packages/python-client/README.md new file mode 100644 index 000000000..00e5c8b07 --- /dev/null +++ b/packages/python-client/README.md @@ -0,0 +1,48 @@ +# python-client + +Connect to a Scrypted server from Python and use the same SDK objects +(`systemManager`, `deviceManager`, `mediaManager`) that Python plugins see. + +The modules in this directory are symlinks into `../../server/python` and +`../../sdk/types` so the client and the plugin runtime share one +implementation. + +## Usage + +```bash +pip install -r requirements.txt +``` + +```python +import asyncio + +from scrypted_client import connect_scrypted_client + + +async def main(loop): + transport, sdk = await connect_scrypted_client( + loop, "https://localhost:10443", "username", "password" + ) + for id in sdk.systemManager.getSystemState(): + print(sdk.systemManager.getDeviceById(id).name) + await transport.close() + + +loop = asyncio.new_event_loop() +loop.run_until_complete(main(loop)) +``` + +`test.py` is a runnable version of the above; point it at a server with the +`SCRYPTED_BASE_URL`, `SCRYPTED_USERNAME`, and `SCRYPTED_PASSWORD` environment +variables. + +By default the login request and the engine.io connection skip TLS +certificate verification, since Scrypted servers use self-signed certificates +out of the box. To control TLS (or connection pooling), pass your own +`login_session` and a pre-built `EioRpcTransport(loop, http_session=...)`. + +Session ownership: a `login_session` you pass is only borrowed for the login +request and is never closed. An `http_session` passed to `EioRpcTransport` +becomes owned by the transport — `transport.close()` closes it along with the +engine.io connection and background tasks — so don't share that session with +anything else. diff --git a/packages/python-client/cluster_setup.py b/packages/python-client/cluster_setup.py new file mode 120000 index 000000000..d80801b0e --- /dev/null +++ b/packages/python-client/cluster_setup.py @@ -0,0 +1 @@ +../../server/python/cluster_setup.py \ No newline at end of file diff --git a/packages/python-client/scrypted_client.py b/packages/python-client/scrypted_client.py new file mode 100644 index 000000000..e0d9d1655 --- /dev/null +++ b/packages/python-client/scrypted_client.py @@ -0,0 +1,249 @@ +"""Client library for connecting to a Scrypted server over engine.io RPC. + +Provides :class:`EioRpcTransport` and :func:`connect_scrypted_client`, the +importable equivalent of the bootstrap that previously only existed inline in +test.py. Typical usage:: + + transport, sdk = await connect_scrypted_client( + asyncio.get_event_loop(), "https://localhost:10443", username, password + ) + for id in sdk.systemManager.getSystemState(): + print(sdk.systemManager.getDeviceById(id).name) + await transport.close() + +Must be called from a running event loop; the transport schedules its send +loop on construction. +""" +from __future__ import annotations + +import asyncio +import contextlib +import logging + +import aiohttp +import engineio + +import plugin_remote +import rpc_reader +from cluster_setup import ClusterSetup +from plugin_remote import DeviceManager, MediaManager, SystemManager +from scrypted_python.scrypted_sdk import ScryptedStatic + +logger = logging.getLogger(__name__) + +DEFAULT_PLUGIN_ID = "@scrypted/core" +DEFAULT_CONNECT_TIMEOUT = 30 + + +class ScryptedConnectionError(Exception): + """Raised when the Scrypted engine.io connection cannot be established.""" + + +class EioRpcTransport(rpc_reader.RpcTransport): + """RpcTransport over an engine.io connection.""" + + def __init__( + self, + loop: asyncio.AbstractEventLoop, + http_session: aiohttp.ClientSession | None = None, + ) -> None: + super().__init__() + # When a session is provided, engineio must defer to its connector + # (ssl_verify=True skips engineio's own ssl.create_default_context, + # a blocking call some frameworks forbid on the event loop). Only + # fall back to engineio's ssl_verify=False when we have no session; + # Scrypted servers use self-signed certificates by default. engineio + # never closes externally provided sessions, so the transport takes + # ownership: close() closes the session. Callers that need to keep + # the session alive should not share it with the transport. + self._http_session = http_session + self.eio = engineio.AsyncClient( + http_session=http_session, ssl_verify=http_session is not None + ) + self.loop = loop + self.write_error: Exception | None = None + self.read_queue: asyncio.Queue = asyncio.Queue() + self.write_queue: asyncio.Queue = asyncio.Queue() + self._send_task: asyncio.Task | None = None + # Set by connect_scrypted_client() so close() can tear it down. + self.read_task: asyncio.Task | None = None + + @self.eio.on("message") + def on_message(data): + self.read_queue.put_nowait(data) + + self._send_task = loop.create_task(self.send_loop()) + + async def read(self): + return await self.read_queue.get() + + async def send_loop(self): + while True: + data = await self.write_queue.get() + try: + await self.eio.send(data) + except Exception as e: + # Any send failure kills the link; surface it on next write. + self.write_error = e + break + + def writeBuffer(self, buffer, reject): + if self.write_error: + if reject: + reject(self.write_error) + return + self.write_queue.put_nowait(buffer) + + def writeSerialized(self, j, reject): + # engineio json-encodes dict payloads on send and json-decodes text + # frames on receive, so messages cross the wire as engine.io JSON + # packets and arrive at readLoop() already deserialized to dicts. + return self.writeBuffer(j, reject) + + async def close(self) -> None: + """Tear the transport down, cancelling background tasks.""" + tasks = [task for task in (self._send_task, self.read_task) if task] + self._send_task = None + self.read_task = None + for task in tasks: + task.cancel() + try: + await self.eio.disconnect() + except Exception: + logger.debug("Error disconnecting engine.io client", exc_info=True) + for task in tasks: + with contextlib.suppress(asyncio.CancelledError, Exception): + await task + if self._http_session: + await self._http_session.close() + self._http_session = None + + +async def connect_scrypted_client( + loop: asyncio.AbstractEventLoop, + base_url: str, + username: str, + password: str, + plugin_id: str = DEFAULT_PLUGIN_ID, + login_session: aiohttp.ClientSession | None = None, + transport: EioRpcTransport | None = None, + timeout: float = DEFAULT_CONNECT_TIMEOUT, +) -> tuple[EioRpcTransport, ScryptedStatic]: + """Login and establish the engine.io RPC session. + + Returns ``(transport, sdk)``. The caller owns the transport and must + ``await transport.close()`` when done. If ``login_session`` is omitted, a + temporary session with certificate verification disabled is used for the + login request only; pass a configured session to control TLS behavior. + Likewise pass a pre-built :class:`EioRpcTransport` to control the + engine.io connection's session (the transport takes ownership of that + session and closes it in ``close()``). ``timeout`` bounds each phase of + the connection: the login request, the engine.io connect, and the wait + for initial system state. Raises :class:`ScryptedConnectionError` on any + failure. + """ + owns_login_session = login_session is None + session = login_session or aiohttp.ClientSession( + connector=aiohttp.TCPConnector(ssl=False) + ) + try: + async with session.post( + f"{base_url}/login", + json={"username": username, "password": password}, + raise_for_status=True, + timeout=aiohttp.ClientTimeout(total=timeout), + ) as response: + login_response = await response.json() + except (aiohttp.ClientError, asyncio.TimeoutError) as err: + raise ScryptedConnectionError(f"Login to {base_url} failed: {err}") from err + finally: + if owns_login_session: + await session.close() + + if "authorization" not in login_response: + raise ScryptedConnectionError( + f"Login to {base_url} did not return an authorization header " + f"(response keys: {sorted(login_response)})" + ) + + if transport is None: + transport = EioRpcTransport(loop) + try: + await asyncio.wait_for( + transport.eio.connect( + base_url, + headers={"Authorization": login_response["authorization"]}, + engineio_path=f"/endpoint/{plugin_id}/engine.io/api/", + ), + timeout, + ) + except (engineio.exceptions.ConnectionError, asyncio.TimeoutError) as err: + await transport.close() + raise ScryptedConnectionError(f"engine.io connect failed: {err}") from err + + ret: asyncio.Future[ScryptedStatic] = loop.create_future() + peer, peer_read_loop = await rpc_reader.prepare_peer_readloop(loop, transport) + peer.params["print"] = logger.debug + + def get_remote(api, plugin_id_, host_info): + cluster_setup = ClusterSetup(loop, peer) + remote = plugin_remote.PluginRemote( + cluster_setup, api, plugin_id_, host_info, loop + ) + wrapped = remote.setSystemState + + async def remote_set_system_state(system_state): + await wrapped(system_state) + + async def resolve(): + if ret.done(): + return + sdk = ScryptedStatic() + sdk.api = api + sdk.remote = remote + system_manager = SystemManager(api, remote.systemState) + sdk.systemManager = system_manager + sdk.deviceManager = DeviceManager(remote.nativeIds, system_manager) + sdk.mediaManager = MediaManager(await api.getMediaManager()) + if not ret.done(): + # PluginRemote.notify dispatches events to the attached + # systemManager's registry (same wiring as loadZip); this + # is what makes systemManager.listen() callbacks fire. + remote.systemManager = system_manager + ret.set_result(sdk) + + loop.create_task(resolve()) + + remote.setSystemState = remote_set_system_state + return remote + + peer.params["getRemote"] = get_remote + read_task = loop.create_task(peer_read_loop()) + transport.read_task = read_task + + def on_read_loop_done(task: asyncio.Task) -> None: + # Surface a dead link immediately instead of waiting for the timeout. + if ret.done() or task.cancelled(): + return + err = task.exception() + ret.set_exception( + ScryptedConnectionError( + f"Connection to {base_url} lost during handshake: {err}" + if err + else f"Connection to {base_url} closed during handshake" + ) + ) + + read_task.add_done_callback(on_read_loop_done) + + try: + sdk = await asyncio.wait_for(asyncio.shield(ret), timeout) + except asyncio.TimeoutError as err: + await transport.close() + raise ScryptedConnectionError( + f"Timed out waiting for Scrypted system state from {base_url}" + ) from err + except ScryptedConnectionError: + await transport.close() + raise + return transport, sdk diff --git a/packages/python-client/test.py b/packages/python-client/test.py index 069fc009e..3be0b1bf7 100644 --- a/packages/python-client/test.py +++ b/packages/python-client/test.py @@ -2,137 +2,15 @@ from __future__ import annotations import asyncio import os -from contextlib import nullcontext -import aiohttp -import engineio - -import plugin_remote -import rpc_reader -from plugin_remote import DeviceManager, MediaManager, SystemManager -from scrypted_python.scrypted_sdk import ScryptedInterface, ScryptedStatic +from scrypted_client import connect_scrypted_client +from scrypted_python.scrypted_sdk import ScryptedInterface -class EioRpcTransport(rpc_reader.RpcTransport): - def __init__(self, loop: asyncio.AbstractEventLoop): - super().__init__() - self.eio = engineio.AsyncClient(ssl_verify=False) - self.loop = loop - self.write_error: Exception = None - self.read_queue = asyncio.Queue() - self.write_queue = asyncio.Queue() - - @self.eio.on("message") - def on_message(data): - self.read_queue.put_nowait(data) - - asyncio.run_coroutine_threadsafe(self.send_loop(), self.loop) - - async def read(self): - return await self.read_queue.get() - - async def send_loop(self): - while True: - data = await self.write_queue.get() - try: - await self.eio.send(data) - except Exception as e: - self.write_error = e - self.write_queue = None - break - - def writeBuffer(self, buffer, reject): - async def send(): - try: - if self.write_error: - raise self.write_error - self.write_queue.put_nowait(buffer) - except Exception as e: - reject(e) - - asyncio.run_coroutine_threadsafe(send(), self.loop) - - def writeSerialized(self, json, reject): - return self.writeBuffer(json, reject) - - -async def connect_scrypted_client( - transport: EioRpcTransport, - base_url: str, - username: str, - password: str, - plugin_id: str = "@scrypted/core", - session: aiohttp.ClientSession | None = None, -) -> ScryptedStatic: - login_url = f"{base_url}/login" - login_body = { - "username": username, - "password": password, - } - - if session: - cm = nullcontext(session) - else: - cm = aiohttp.ClientSession() - - async with cm as _session: - async with _session.post( - login_url, verify_ssl=False, json=login_body - ) as response: - login_response = await response.json() - - headers = {"Authorization": login_response["authorization"]} - - await transport.eio.connect( - base_url, - headers=headers, - engineio_path=f"/endpoint/{plugin_id}/engine.io/api/", - ) - - ret = asyncio.Future[ScryptedStatic](loop=transport.loop) - peer, peerReadLoop = await rpc_reader.prepare_peer_readloop( - transport.loop, transport - ) - peer.params["print"] = print - - def callback(api, pluginId, hostInfo): - cluster_setup = plugin_remote.ClusterSetup(transport.loop, peer) - remote = plugin_remote.PluginRemote( - cluster_setup, api, pluginId, hostInfo, transport.loop - ) - wrapped = remote.setSystemState - - async def remoteSetSystemState(systemState): - await wrapped(systemState) - - async def resolve(): - sdk = ScryptedStatic() - sdk.api = api - sdk.remote = remote - sdk.systemManager = SystemManager(api, remote.systemState) - sdk.deviceManager = DeviceManager( - remote.nativeIds, sdk.systemManager - ) - sdk.mediaManager = MediaManager(await api.getMediaManager()) - ret.set_result(sdk) - - asyncio.run_coroutine_threadsafe(resolve(), transport.loop) - - remote.setSystemState = remoteSetSystemState - return remote - - peer.params["getRemote"] = callback - asyncio.run_coroutine_threadsafe(peerReadLoop(), transport.loop) - - sdk = await ret - return sdk - - -async def main(): - transport = EioRpcTransport(asyncio.get_event_loop()) - sdk = await connect_scrypted_client( - transport, - "https://localhost:10443", +async def main(loop: asyncio.AbstractEventLoop): + transport, sdk = await connect_scrypted_client( + loop, + os.environ.get("SCRYPTED_BASE_URL", "https://localhost:10443"), os.environ["SCRYPTED_USERNAME"], os.environ["SCRYPTED_PASSWORD"], ) @@ -143,10 +21,12 @@ async def main(): if ScryptedInterface.OnOff.value in device.interfaces: print(f"OnOff: device is {device.on}") - await transport.eio.disconnect() - os._exit(0) + await transport.close() -loop = asyncio.new_event_loop() -asyncio.run_coroutine_threadsafe(main(), loop) -loop.run_forever() +if __name__ == "__main__": + loop = asyncio.new_event_loop() + try: + loop.run_until_complete(main(loop)) + finally: + loop.close() diff --git a/server/python/plugin_remote.py b/server/python/plugin_remote.py index 98a32c6aa..10fddb1ac 100644 --- a/server/python/plugin_remote.py +++ b/server/python/plugin_remote.py @@ -20,14 +20,10 @@ from io import StringIO from pathlib import Path from typing import Any, Callable, Coroutine, Optional, Set, Tuple, TypedDict -import cluster_labels -import plugin_console -import plugin_volume as pv import rpc import rpc_reader import scrypted_python.scrypted_sdk.types from cluster_setup import ClusterSetup -from plugin_pip import install_with_pip, need_requirements, remove_pip_dirs from scrypted_python.scrypted_sdk import PluginFork, ScryptedStatic from scrypted_python.scrypted_sdk.types import (Device, DeviceManifest, EventDetails, @@ -723,6 +719,13 @@ class PluginRemote: raise async def loadZipWrapped(self, packageJson, zipAPI: Any, zipOptions: dict): + # plugin host modules, unused (and not installed) when this module is + # consumed by a client via the scrypted-sdk package + import cluster_labels + import plugin_console + import plugin_volume as pv + from plugin_pip import install_with_pip, need_requirements, remove_pip_dirs + await self.clusterSetup.initializeCluster(zipOptions) sdk = ScryptedStatic()