diff --git a/rpc/src/core/client.ts b/rpc/src/core/client.ts index 12d5a55..dd40f06 100644 --- a/rpc/src/core/client.ts +++ b/rpc/src/core/client.ts @@ -20,19 +20,13 @@ import { Wrapper } from './wrapper' export class RPCInstanceClient extends Wrapper { private readonly _emitterServer: Emitter - private readonly _pendingServer: Emitter private readonly _emitterWeb: Emitter - private readonly _pendingWeb: Emitter - private readonly _pendingWebToServer: Emitter constructor(props: RPCConfig<'client'>) { super(props) this._emitterServer = new Emitter() - this._pendingServer = new Emitter() this._emitterWeb = new Emitter() - this._pendingWeb = new Emitter() - this._pendingWebToServer = new Emitter() this.console.log('[RPC] Initialized Client') @@ -93,17 +87,8 @@ export class RPCInstanceClient extends Wrapper { } } if (payload.type === 'response') { - if (payload.calledTo === 'client') { - await this._pendingServer.emit( - payload.uuid, - ...(payload.data && payload.data.length > 0 ? payload.data : []), - ) - } - if (payload.calledTo === 'webview') { - await this._pendingWebToServer.emit( - payload.uuid, - ...(payload.data && payload.data.length > 0 ? payload.data : []), - ) + if (payload.calledTo === 'client' || payload.calledTo === 'webview') { + this.resolvePending(payload) } } } @@ -126,18 +111,13 @@ export class RPCInstanceClient extends Wrapper { payload.player = GetPlayerServerId(PlayerId()) emitNet(RPCEvents.LISTENER_WEB, stringify(payload)) - return new Promise(res => { - this._pendingWebToServer.once(payload.uuid, res) - }) + return this._pending.wait(payload.uuid) } } if (payload.type === 'response') { if (payload.calledTo === 'client') { - await this._pendingWeb.emit( - payload.uuid, - ...(payload.data && payload.data.length > 0 ? payload.data : []), - ) + this.resolvePending(payload) return { status: 'ok' } } @@ -202,9 +182,7 @@ export class RPCInstanceClient extends Wrapper { emitNet(RPCEvents.LISTENER_CLIENT, stringify(payload)) - return new Promise>(res => { - this._pendingServer.once(payload.uuid, res) - }) + return this._pending.wait>(payload.uuid) } // ===== WEBVIEW ===== @@ -261,9 +239,7 @@ export class RPCInstanceClient extends Wrapper { data: payload, }) - return new Promise>(res => { - this._pendingWeb.once(payload.uuid, res) - }) + return this._pending.wait>(payload.uuid) } // ===== SELF ===== diff --git a/rpc/src/core/server.ts b/rpc/src/core/server.ts index 8cb72fc..0340a60 100644 --- a/rpc/src/core/server.ts +++ b/rpc/src/core/server.ts @@ -15,17 +15,13 @@ import { Wrapper } from './wrapper' export class RPCInstanceServer extends Wrapper { private readonly _emitterClient: Emitter - private readonly _pendingClient: Emitter private readonly _emitterWeb: Emitter - private readonly _pendingWeb: Emitter constructor(props: RPCConfig<'server'>) { super(props) this._emitterClient = new Emitter() - this._pendingClient = new Emitter() this._emitterWeb = new Emitter() - this._pendingWeb = new Emitter() this.console.log('[RPC] Initialized Server') @@ -78,10 +74,7 @@ export class RPCInstanceServer extends Wrapper { emitNet(RPCEvents.LISTENER_SERVER, response.player, stringify(response)) } if (payload.type === 'response') { - await this._pendingClient.emit( - payload.uuid, - ...(payload.data && payload.data.length > 0 ? payload.data : []), - ) + this.resolvePending(payload) } } } @@ -129,10 +122,7 @@ export class RPCInstanceServer extends Wrapper { emitNet(RPCEvents.LISTENER_SERVER, response.player, stringify(response)) } if (payload.type === 'response') { - await this._pendingWeb.emit( - payload.uuid, - ...(payload.data && payload.data.length > 0 ? payload.data : []), - ) + this.resolvePending(payload) } } } @@ -193,9 +183,7 @@ export class RPCInstanceServer extends Wrapper { emitNet(RPCEvents.LISTENER_SERVER, player, stringify(payload)) - return new Promise>(res => { - this._pendingClient.once(payload.uuid, res) - }) + return this._pending.wait>(payload.uuid) } public async emitClientEveryone< @@ -272,9 +260,7 @@ export class RPCInstanceServer extends Wrapper { emitNet(RPCEvents.LISTENER_SERVER, player, stringify(payload)) - return new Promise>(res => { - this._pendingWeb.once(payload.uuid, res) - }) + return this._pending.wait>(payload.uuid) } // ===== SELF ===== diff --git a/rpc/src/core/webview.ts b/rpc/src/core/webview.ts index 4f0a1d5..d269ef4 100644 --- a/rpc/src/core/webview.ts +++ b/rpc/src/core/webview.ts @@ -144,7 +144,7 @@ export class RPCInstanceWebview extends Wrapper { type: 'event', } - return await this._createHttpClientRequest>(payload) + return this._request>(payload) } // ===== SERVER ===== @@ -196,7 +196,7 @@ export class RPCInstanceWebview extends Wrapper { type: 'event', } - return await this._createHttpClientRequest>(payload) + return this._request>(payload) } // ===== SELF ===== @@ -264,6 +264,16 @@ export class RPCInstanceWebview extends Wrapper { // ===== UTILS ===== + /** Sends an event to the client and waits for its response (with timeout) */ + private _request(payload: RPCState): Promise { + const response = this._pending.wait(payload.uuid) + this._createHttpClientRequest(payload).then( + data => this._pending.resolve(payload.uuid, data), + (error: Error) => this._pending.reject(payload.uuid, error), + ) + return response + } + private async _createHttpClientRequest( data: RPCStateRaw | RPCState, ): Promise { diff --git a/rpc/src/core/wrapper.ts b/rpc/src/core/wrapper.ts index 86f3061..a36d4f6 100644 --- a/rpc/src/core/wrapper.ts +++ b/rpc/src/core/wrapper.ts @@ -1,5 +1,6 @@ import { Emitter } from '../utils/emitter' import { parse } from '../utils/funcs' +import { Pending } from '../utils/pending' import { type RPCConfig, type RPCEnvironment, @@ -11,16 +12,28 @@ import { export class Wrapper { protected env: RPCEnvironment protected _emitterLocal: Emitter + protected _pending: Pending protected debug: boolean protected console: Console constructor(cfg: RPCConfig) { this.env = cfg.env this._emitterLocal = new Emitter() + this._pending = new Pending(cfg.timeout ?? 5000) this.debug = cfg.debug ?? false this.console = console } + /** Settles the call waiting for this response; ignores late or unexpected ones */ + protected resolvePending(payload: RPCState): void { + const found = this._pending.resolve(payload.uuid, payload.data?.[0]) + if (!found && this.debug) { + this.console.log( + `[RPC]:ignored response ${payload.event} ${payload.uuid} (no pending call, possibly timed out)`, + ) + } + } + protected verifyEvent(state: Emitter, data: RPCStateRaw | RPCState) { const rpcData = typeof data === 'string' ? parse(data) : data diff --git a/rpc/src/utils/pending.ts b/rpc/src/utils/pending.ts new file mode 100644 index 0000000..d1af6f5 --- /dev/null +++ b/rpc/src/utils/pending.ts @@ -0,0 +1,56 @@ +import { RPCErrors } from './types' + +type Call = { + resolve: (data: unknown) => void + reject: (error: Error) => void + timer: ReturnType | undefined +} + +/** Calls waiting for a response, keyed by payload uuid. */ +export class Pending { + private _calls = new Map() + + /** @param timeout - ms before a call rejects, `0` or less disables it */ + constructor(private readonly _timeout: number) {} + + public wait(uuid: string): Promise { + return new Promise((resolve, reject) => { + const timer = + this._timeout > 0 + ? setTimeout( + () => this.reject(uuid, new Error(RPCErrors.TIMEOUT)), + this._timeout, + ) + : undefined + + this._calls.set(uuid, { + resolve: resolve as (data: unknown) => void, + reject, + timer, + }) + }) + } + + /** @returns `false` if no call waits for this uuid (late or unexpected response) */ + public resolve(uuid: string, data: unknown): boolean { + const call = this._take(uuid) + call?.resolve(data) + return call !== undefined + } + + /** @returns `false` if no call waits for this uuid */ + public reject(uuid: string, error: Error): boolean { + const call = this._take(uuid) + call?.reject(error) + return call !== undefined + } + + private _take(uuid: string): Call | undefined { + const call = this._calls.get(uuid) + if (call) { + this._calls.delete(uuid) + clearTimeout(call.timer) + } + return call + } +} diff --git a/rpc/src/utils/types.ts b/rpc/src/utils/types.ts index e08ad43..23a23b0 100644 --- a/rpc/src/utils/types.ts +++ b/rpc/src/utils/types.ts @@ -24,6 +24,13 @@ export type RPCEnvironmentResolved = export type RPCConfig = { env: T debug?: boolean + /** + * Milliseconds to wait for a response before the call rejects with + * `RPCErrors.TIMEOUT`. `0` disables the timeout. + * + * @defaultValue 5000 + */ + timeout?: number } /** **Internal** */ @@ -85,6 +92,7 @@ export enum RPCErrors { NO_PLAYER = 'No player (failed to resolve from local index)', UNKNOWN_NATIVE = 'Unknown native event (if you are sure this exists - use native handler)', UNKNOWN_ENVIRONMENT = 'Unknown environment (must be either "server", "client" or "webview")', + TIMEOUT = 'Timed out waiting for response', } /**