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 <koushd@gmail.com>
This commit is contained in:
Brett Jia
2023-11-08 12:22:07 -05:00
committed by GitHub
parent d7a417c984
commit e49f26b410
7 changed files with 269 additions and 15 deletions

27
packages/client/.vscode/launch.json vendored Normal file
View File

@@ -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"
}
]
}

View File

@@ -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<VideoCamera & Camera>("Office");
const libav = sdk.systemManager.getDeviceByName<VideoFrameGenerator>("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();

View File

@@ -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<number, Promise<{ clusterPeer: RpcPeer, clusterSecret: string }>>();
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', () => {

View File

@@ -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<any>;
/*
* 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));
};

View File

@@ -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<any>;
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);

View File

@@ -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,

View File

@@ -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<HttpPluginData> {
})
},
});
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<HttpPluginData> {
}
});
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<HttpPluginData> {
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<HttpPluginData> {
const ret = await this.getPluginForEndpoint(endpoint);
if (req.url.indexOf('/engine.io/api') !== -1)