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()