Something went wrong. Try again.
collection of js libraries
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299import z from "zod";import type { RpcContract, RpcContractDefinition, RpcProcedureDefinition } from "./contract.ts";import type { RpcTransport } from "./transports/rpc-transport.ts";import type { Message } from "./messages.ts";import { PendingResponses } from "./pending-responses.ts";import { v7 as uuidv7 } from "uuid";
type TimeoutId = ReturnType<typeof setTimeout>;
export type RpcPeerOptions = { transport: RpcTransport; localContract: RpcContract<any>; remoteContract: RpcContract<any>;};
type MaybeAsyncGenerator<T, TReturn, TNext> = | AsyncGenerator<T, TReturn, TNext> | Generator<T, TReturn, TNext>;type MaybePromise<T> = T | Promise<T>;
type ProcedureResult<T_Definition extends RpcProcedureDefinition> = T_Definition["generator"] extends true ? MaybeAsyncGenerator<z.infer<T_Definition["output"]>, void, unknown> : MaybePromise<z.infer<T_Definition["output"]>>;
type LocalHandlers<T_Definition extends RpcContractDefinition> = { [Method in keyof T_Definition]: ( params: z.infer<T_Definition[Method]["input"]>, ) => ProcedureResult<T_Definition[Method]>;};
type RemoteCallReturn<T_Definition extends RpcProcedureDefinition> = T_Definition["generator"] extends true ? AsyncGenerator<z.infer<T_Definition["output"]>, void, unknown> : Awaited<z.infer<T_Definition["output"]>>;
const STREAM_IDLE_TIMEOUT_MS = 30000;
type StreamEntry = { streamId: string; generator: MaybeAsyncGenerator<any, void, unknown>; timeoutId?: TimeoutId;};
export class RpcPeer<T_Options extends RpcPeerOptions> { private _options: T_Options;
private _handlers: LocalHandlers<T_Options["localContract"]["definition"]>;
private _pendingResponses: PendingResponses<Message>;
private _streams: Map<string, StreamEntry>;
constructor( options: T_Options, handlers: LocalHandlers<T_Options["localContract"]["definition"]>, ) { this._options = options; this._handlers = handlers;
this._pendingResponses = new PendingResponses(); this._streams = new Map(); this._options.transport.onMessage((message) => { this._handleIncomingMessage(message); }); }
async call<T_Method extends keyof T_Options["remoteContract"]["definition"]>( method: T_Method, params: z.infer<T_Options["remoteContract"]["definition"][T_Method]["input"]>, ): Promise<RemoteCallReturn<T_Options["remoteContract"]["definition"][T_Method]>> { const msgId = uuidv7(); const response = await this._sendMessageAndWaitForResponse({ type: "callMethod", method: method as string, value: params, msgId, });
if (response.type === "callResponse") { return response.value as RemoteCallReturn< T_Options["remoteContract"]["definition"][T_Method] >; } else if (response.type === "callResponseStream") { return this._createStreamGenerator(response.streamId, response.senderId) as any; } else { throw new Error(`Unexpected message type: ${response.type}`); } }
private async *_createStreamGenerator(streamId: string, senderId: string | undefined) { while (true) { const msgId = uuidv7(); const chunkResponse = await this._sendMessageAndWaitForResponse({ type: "requestStreamChunk", streamId, msgId, value: undefined, senderId, });
if (chunkResponse.type === "streamChunkValue") { yield chunkResponse.value; } else if (chunkResponse.type === "streamChunkDone") { return; } else if (chunkResponse.type === "streamChunkError") { throw new Error(`Stream error from remote: ${chunkResponse.error}`); } else { throw new Error(`Unexpected message type: ${chunkResponse.type}`); } } }
private _refreshStreamTimeout(streamId: string) { const stream = this._streams.get(streamId); if (!stream) return;
if (stream.timeoutId) { clearTimeout(stream.timeoutId); }
stream.timeoutId = setTimeout(() => { this._streams.delete(streamId); }, STREAM_IDLE_TIMEOUT_MS); }
private _deleteStream(streamId: string) { const stream = this._streams.get(streamId); if (!stream) return;
if (stream.timeoutId) { clearTimeout(stream.timeoutId); } this._streams.delete(streamId); }
private async _handleIncomingMessage(message: Message) { if ("prevMsgId" in message) { this._pendingResponses.resolve(message.prevMsgId, message); }
if (message.type === "ping") { this._options.transport.sendMessage({ type: "pong", msgId: uuidv7(), prevMsgId: message.msgId, }); return; } else if (message.type === "pong") { return; } else if (message.type === "callMethod") { const definition = this._options.localContract.definition[ message.method as keyof T_Options["localContract"]["definition"] ]; if (!definition) { this._sendMessage({ type: "callResponse", value: undefined, msgId: uuidv7(), prevMsgId: message.msgId, senderId: message.senderId, error: `No method named ${message.method}`, }); return; }
const handler = this._handlers[message.method as keyof T_Options["localContract"]["definition"]]; if (!handler) { this._sendMessage({ type: "callResponse", value: undefined, msgId: uuidv7(), prevMsgId: message.msgId, senderId: message.senderId, error: `No handler for method ${message.method}`, }); return; }
try { const result = await handler(message.value);
if (isGenerator(result) || isAsyncGenerator(result)) { const streamId = uuidv7(); this._streams.set(streamId, { streamId, generator: result }); this._refreshStreamTimeout(streamId); this._sendMessage({ type: "callResponseStream", msgId: uuidv7(), prevMsgId: message.msgId, senderId: message.senderId, streamId, }); } else { this._sendMessage({ type: "callResponse", msgId: uuidv7(), prevMsgId: message.msgId, senderId: message.senderId, value: result, }); } } catch (err) { this._sendMessage({ type: "callResponse", msgId: uuidv7(), value: undefined, prevMsgId: message.msgId, senderId: message.senderId, error: err instanceof Error ? err.message : String(err), }); }
return; } else if (message.type === "requestStreamChunk") { const stream = this._streams.get(message.streamId); if (!stream) { this._sendMessage({ type: "streamChunkError", msgId: uuidv7(), prevMsgId: message.msgId, streamId: message.streamId, error: `No stream with id ${message.streamId}`, senderId: message.senderId, }); return; }
try { const { value, done } = await stream.generator.next(); if (done) { this._sendMessage({ type: "streamChunkDone", msgId: uuidv7(), prevMsgId: message.msgId, streamId: message.streamId, senderId: message.senderId, }); this._deleteStream(message.streamId); } else { this._refreshStreamTimeout(message.streamId); this._sendMessage({ type: "streamChunkValue", msgId: uuidv7(), prevMsgId: message.msgId, streamId: message.streamId, senderId: message.senderId, value, }); } } catch (err) { this._options.transport.sendMessage({ type: "streamChunkError", msgId: uuidv7(), prevMsgId: message.msgId, streamId: message.streamId, error: err instanceof Error ? err.message : String(err), senderId: message.senderId, }); this._streams.delete(message.streamId); }
return; } else if (message.type === "callResponse") { return; } }
private _sendMessage(message: Message) { this._options.transport.sendMessage(message); }
private async _sendMessageAndWaitForResponse(message: Message): Promise<Message> { this._options.transport.sendMessage(message); return this._pendingResponses.waitFor(message.msgId); }}
function isGenerator(obj: any): obj is Generator { return ( obj && typeof obj.next === "function" && typeof obj.throw === "function" && typeof obj.return === "function" );}
function isAsyncGenerator(obj: any): obj is AsyncGenerator { return ( obj && typeof obj.next === "function" && typeof obj.throw === "function" && typeof obj.return === "function" && obj[Symbol.toStringTag] === "AsyncGenerator" );}