From e49f26b410f86e63b4487e4f8440e930cabfdb38 Mon Sep 17 00:00:00 2001 From: Brett Jia Date: Wed, 8 Nov 2023 12:22:07 -0500 Subject: [PATCH] server, client: connectRPCObject for web api clients (#1166) * initial pass at connectRPCObject proxy * wip: connects to server but fails on pendingResults * fix wrong rpcpeer bug + cleanup serialization * small cleanups * feedback, local object lookups * rpc: fix up additional id gens * feedback * update example to use frame generator --------- Co-authored-by: Koushik Dutta --- packages/client/.vscode/launch.json | 27 +++++ packages/client/examples/connectRPCObject.ts | 34 ++++++ packages/client/src/index.ts | 111 ++++++++++++++++++- server/src/plugin/connect-rpc-object.ts | 30 +++++ server/src/plugin/plugin-remote-worker.ts | 10 +- server/src/rpc.ts | 10 +- server/src/runtime.ts | 62 +++++++++++ 7 files changed, 269 insertions(+), 15 deletions(-) create mode 100644 packages/client/.vscode/launch.json create mode 100644 packages/client/examples/connectRPCObject.ts create mode 100644 server/src/plugin/connect-rpc-object.ts diff --git a/packages/client/.vscode/launch.json b/packages/client/.vscode/launch.json new file mode 100644 index 000000000..9d95af636 --- /dev/null +++ b/packages/client/.vscode/launch.json @@ -0,0 +1,27 @@ +{ + // Use IntelliSense to learn about possible attributes. + // Hover to view descriptions of existing attributes. + // For more information, visit: https://go.microsoft.com/fwlink/?linkid=830387 + "version": "0.2.0", + "configurations": [ + { + "name": "ts-node", + "type": "node", + "request": "launch", + "args": [ + "${relativeFile}" + ], + "runtimeArgs": [ + "-r", + "ts-node/register" + ], + "env": { + "SCRYPTED_USERNAME": "koush", + "SCRYPTED_PASSWORD": "k9copUSA", + }, + "cwd": "${workspaceRoot}", + "protocol": "inspector", + "internalConsoleOptions": "openOnSessionStart" + } + ] +} \ No newline at end of file diff --git a/packages/client/examples/connectRPCObject.ts b/packages/client/examples/connectRPCObject.ts new file mode 100644 index 000000000..b9d738bfe --- /dev/null +++ b/packages/client/examples/connectRPCObject.ts @@ -0,0 +1,34 @@ +import { Camera, VideoCamera, VideoFrameGenerator } from '@scrypted/types'; +import { connectScryptedClient } from '../dist/packages/client/src'; + +import https from 'https'; + +const httpsAgent = new https.Agent({ + rejectUnauthorized: false, +}) + +async function example() { + const sdk = await connectScryptedClient({ + baseUrl: 'https://localhost:10443', + pluginId: "@scrypted/core", + username: process.env.SCRYPTED_USERNAME || 'admin', + password: process.env.SCRYPTED_PASSWORD || 'swordfish', + axiosConfig: { + httpsAgent, + } + }); + console.log('server version', sdk.serverVersion); + + const office = sdk.systemManager.getDeviceByName("Office"); + const libav = sdk.systemManager.getDeviceByName("Libav"); + const mo = await office.getVideoStream(); + + const generator = await libav.generateVideoFrames(mo); + const remote = await sdk.connectRPCObject(generator); + + for await (const frame of remote) { + console.log(frame); + } +} + +example(); diff --git a/packages/client/src/index.ts b/packages/client/src/index.ts index 633dcbb64..18bedabc3 100644 --- a/packages/client/src/index.ts +++ b/packages/client/src/index.ts @@ -1,3 +1,4 @@ +import crypto from 'crypto'; import { MediaObjectOptions, RTCConnectionManagement, RTCSignalingSession, ScryptedStatic } from "@scrypted/types"; import axios, { AxiosRequestConfig, AxiosRequestHeaders } from 'axios'; import * as eio from 'engine.io-client'; @@ -9,11 +10,14 @@ import { DataChannelDebouncer } from "../../../plugins/webrtc/src/datachannel-de import type { IOSocket } from '../../../server/src/io'; import { MediaObject } from '../../../server/src/plugin/mediaobject'; import { attachPluginRemote } from '../../../server/src/plugin/plugin-remote'; +import type { ClusterObject, ConnectRPCObject } from '../../../server/src/plugin/connect-rpc-object'; import { RpcPeer } from '../../../server/src/rpc'; import { createRpcDuplexSerializer, createRpcSerializer } from '../../../server/src/rpc-serializer'; import packageJson from '../package.json'; import { isIPAddress } from "./ip"; +const sourcePeerId = RpcPeer.generateId(); + type IOClientSocket = eio.Socket & IOSocket; function once(socket: IOClientSocket, event: 'open' | 'message') { @@ -707,6 +711,110 @@ export async function connectScryptedClient(options: ScryptedClientOptions): Pro .map(id => systemManager.getDeviceById(id)) .find(device => device.pluginId === '@scrypted/core' && device.nativeId === `user:${username}`); + const clusterPeers = new Map>(); + const ensureClusterPeer = (port: number) => { + let clusterPeerPromise = clusterPeers.get(port); + if (!clusterPeerPromise) { + clusterPeerPromise = (async () => { + const eioPath = 'engine.io/connectRPCObject'; + const eioEndpoint = baseUrl ? new URL(eioPath, baseUrl).pathname : '/' + eioPath; + const clusterPeerOptions = { + path: eioEndpoint, + query: { + cacehBust, + port, + }, + withCredentials: true, + extraHeaders, + rejectUnauthorized: false, + transports: options?.transports, + }; + + const clusterPeerSocket = new eio.Socket(explicitBaseUrl, clusterPeerOptions); + let peerReady = false; + clusterPeerSocket.on('close', () => { + clusterPeers.delete(port); + if (!peerReady) { + throw new Error("peer disconnected before setup completed"); + } + }); + + try { + const clusterSecretPromise = once(clusterPeerSocket, 'message'); + + await once(clusterPeerSocket, 'open'); + + const clusterSecret = await clusterSecretPromise as any as string; + + const serializer = createRpcDuplexSerializer({ + write: data => clusterPeerSocket.send(data), + }); + clusterPeerSocket.on('message', data => serializer.onData(data)); + + const clusterPeer = new RpcPeer(clientName || 'engine.io-client', "cluster-proxy", (message, reject, serializationContext) => { + try { + serializer.sendMessage(message, reject, serializationContext); + } + catch (e) { + reject?.(e); + } + }); + serializer.setupRpcPeer(clusterPeer); + clusterPeer.tags.localPort = sourcePeerId; + peerReady = true; + return { clusterPeer, clusterSecret }; + } + catch (e) { + console.error('failure ipc connect', e); + clusterPeerSocket.close(); + throw e; + } + })(); + clusterPeers.set(port, clusterPeerPromise); + } + return clusterPeerPromise; + }; + + const resolveObject = async (proxyId: string, sourcePeerPort: number) => { + const sourcePeer = (await clusterPeers.get(sourcePeerPort))?.clusterPeer; + if (sourcePeer?.remoteWeakProxies) { + return Object.values(sourcePeer.remoteWeakProxies).find( + v => v.deref()?.__cluster?.proxyId == proxyId + )?.deref(); + } + return null; + } + + const connectRPCObject = async (value: any) => { + const clusterObject: ClusterObject = value?.__cluster; + if (!clusterObject) { + return value; + } + + const { port, proxyId, source } = clusterObject; + + // check if object is already connected + const resolved = await resolveObject(proxyId, port); + if (resolved) { + return resolved; + } + + try { + const clusterPeerPromise = ensureClusterPeer(port); + const { clusterPeer, clusterSecret } = await clusterPeerPromise; + const connectRPCObject: ConnectRPCObject = await clusterPeer.getParam('connectRPCObject'); + const portSecret = crypto.createHash('sha256').update(`${port}${clusterSecret}`).digest().toString('hex'); + const newValue = await connectRPCObject(proxyId, portSecret, source); + if (!newValue) + throw new Error('ipc object not found?'); + return newValue; + } + catch (e) { + console.error('failure ipc', e); + return value; + } + } + const ret: ScryptedClientStatic = { userId: userDevice?.id, serverVersion, @@ -736,7 +844,8 @@ export async function connectScryptedClient(options: ScryptedClientOptions): Pro queryToken, authorization, cloudAddress, - } + }, + connectRPCObject, } socket.on('close', () => { diff --git a/server/src/plugin/connect-rpc-object.ts b/server/src/plugin/connect-rpc-object.ts new file mode 100644 index 000000000..3486af81c --- /dev/null +++ b/server/src/plugin/connect-rpc-object.ts @@ -0,0 +1,30 @@ +import net from "net"; +import { Socket } from "engine.io"; +import { IOSocket } from "../io"; + +export interface ClusterObject { + id: string; + port: number; + proxyId: string; + source: number; +} + +export type ConnectRPCObject = (id: string, secret: string, sourcePeerPort: number) => Promise; + +/* + * Handle incoming connections that will be + * proxied to a connectRPCObject socket. + */ +export function setupConnectRPCObjectProxy(clusterSecret: string, port: number, connection: Socket & IOSocket) { + if (!port) { + throw new Error("invalid port"); + } + + connection.send(clusterSecret); + + const socket = net.connect(port, '127.0.0.1'); + socket.on('close', () => connection.close()); + socket.on('data', data => connection.send(data)); + connection.on('close', () => socket.destroy()); + connection.on('message', message => socket.write(message)); +}; diff --git a/server/src/plugin/plugin-remote-worker.ts b/server/src/plugin/plugin-remote-worker.ts index 19bc8238e..f9fabfa67 100644 --- a/server/src/plugin/plugin-remote-worker.ts +++ b/server/src/plugin/plugin-remote-worker.ts @@ -17,6 +17,7 @@ import { attachPluginRemote, DeviceManagerImpl, PluginReader, setupPluginRemote import { PluginStats, startStatsUpdater } from './plugin-remote-stats'; import { createREPLServer } from './plugin-repl'; import { NodeThreadWorker } from './runtime/node-thread-worker'; +import { ClusterObject, ConnectRPCObject } from './connect-rpc-object'; import crypto from 'crypto'; const { link } = require('linkfs'); @@ -26,15 +27,6 @@ export interface StartPluginRemoteOptions { onClusterPeer(peer: RpcPeer): void; } -interface ClusterObject { - id: string; - port: number; - proxyId: string; - source: number; -} - -type ConnectRPCObject = (id: string, secret: string, sourcePeerPort: number) => Promise; - export function startPluginRemote(mainFilename: string, pluginId: string, peerSend: (message: RpcMessage, reject?: (e: Error) => void, serializationContext?: any) => void, startPluginRemoteOptions?: StartPluginRemoteOptions) { const peer = new RpcPeer('unknown', 'host', peerSend); diff --git a/server/src/rpc.ts b/server/src/rpc.ts index ccce4ccc6..681ab42c6 100644 --- a/server/src/rpc.ts +++ b/server/src/rpc.ts @@ -421,7 +421,7 @@ export class RpcPeer { return !value || (!value[RpcPeer.PROPERTY_JSON_DISABLE_SERIALIZATION] && this.transportSafeArgumentTypes.has(value.constructor?.name)); } - generateId() { + static generateId() { return [...new Array(8)].map(() => RpcPeer.RANDOM_DIGITS.charAt(Math.floor(Math.random() * RpcPeer.RANDOM_DIGITS.length))).join(''); } @@ -430,7 +430,7 @@ export class RpcPeer { return Promise.reject(new RPCResultError(this, 'RpcPeer has been killed (createPendingResult)')); const promise = new Promise((resolve, reject) => { - const id = this.generateId(); + const id = RpcPeer.generateId(); this.pendingResults[id] = { resolve, reject, method }; cb(id, e => reject(new RPCResultError(this, e.message, e))); @@ -626,7 +626,7 @@ export class RpcPeer { let proxiedEntry = this.localProxied.get(value); if (proxiedEntry) { - const __remote_proxy_finalizer_id = this.generateId(); + const __remote_proxy_finalizer_id = RpcPeer.generateId(); proxiedEntry.finalizerId = __remote_proxy_finalizer_id; const ret: RpcRemoteProxyValue = { __remote_proxy_id: proxiedEntry.id, @@ -648,7 +648,7 @@ export class RpcPeer { this.onProxyTypeSerialization.get(__remote_constructor_name)?.(value); - const __remote_proxy_id = this.generateId(); + const __remote_proxy_id = RpcPeer.generateId(); proxiedEntry = { id: __remote_proxy_id, finalizerId: __remote_proxy_id, @@ -839,7 +839,7 @@ export function getEvalSource() { ${RpcPeer} ${startPeriodicGarbageCollection} - + return { startPeriodicGarbageCollection, RpcPeer, diff --git a/server/src/runtime.ts b/server/src/runtime.ts index 7868c71c6..83964da76 100644 --- a/server/src/runtime.ts +++ b/server/src/runtime.ts @@ -34,6 +34,7 @@ import { getPluginVolume } from './plugin/plugin-volume'; import { NodeForkWorker } from './plugin/runtime/node-fork-worker'; import { PythonRuntimeWorker } from './plugin/runtime/python-worker'; import { RuntimeWorker, RuntimeWorkerOptions } from './plugin/runtime/runtime-worker'; +import { setupConnectRPCObjectProxy } from './plugin/connect-rpc-object'; import { getIpAddress, SCRYPTED_INSECURE_PORT, SCRYPTED_SECURE_PORT } from './server-settings'; import { AddressSettings } from './services/addresses'; import { Alerts } from './services/alerts'; @@ -82,6 +83,17 @@ export class ScryptedRuntime extends PluginHttp { }) }, }); + connectRPCObjectIO: IOServer = new io.Server({ + pingTimeout: 120000, + perMessageDeflate: true, + cors: (req, callback) => { + const header = this.getAccessControlAllowOrigin(req.headers); + callback(undefined, { + origin: header, + credentials: true, + }) + }, + }); pluginComponent = new PluginComponent(this); serviceControl = new ServiceControl(this); alerts = new Alerts(this); @@ -145,6 +157,29 @@ export class ScryptedRuntime extends PluginHttp { } }); + app.all('/engine.io/connectRPCObject', (req, res) => { + if (res.locals.aclId) { + res.writeHead(401); + res.end(); + return; + } + if (!req.query.port) { + res.writeHead(404); + res.end(); + return; + } + this.connectRPCObjectHandler(req, res); + }); + + this.connectRPCObjectIO.on('connection', connection => { + try { + const clusterObjectPortHeader = (connection.request as Request).query.port as string; + setupConnectRPCObjectProxy(this.clusterSecret, parseInt(clusterObjectPortHeader), connection); + } catch { + connection.close(); + } + }); + insecure.on('upgrade', (req, socket, upgradeHead) => { (req as any).upgradeHead = upgradeHead; (app as any).handle(req, { @@ -279,6 +314,33 @@ export class ScryptedRuntime extends PluginHttp { this.shellio.handleRequest(req, res); } + async connectRPCObjectHandler(req: Request, res: Response) { + const isUpgrade = isConnectionUpgrade(req.headers); + + const end = (code: number, message: string) => { + if (isUpgrade) { + const socket = res.socket; + socket.write(`HTTP/1.1 ${code} ${message}\r\n` + + '\r\n'); + socket.destroy(); + } + else { + res.status(code); + res.send(message); + } + }; + + if (!res.locals.username) { + end(401, 'Not Authorized'); + return; + } + + if ((req as any).upgradeHead) + this.connectRPCObjectIO.handleUpgrade(req, res.socket, (req as any).upgradeHead) + else + this.connectRPCObjectIO.handleRequest(req, res); + } + async getEndpointPluginData(req: Request, endpoint: string, isUpgrade: boolean, isEngineIOEndpoint: boolean): Promise { const ret = await this.getPluginForEndpoint(endpoint); if (req.url.indexOf('/engine.io/api') !== -1)