import { Emitter } from "@noya-app/emitter"; import { uuid } from "@noya-app/noya-utils"; import { Observable } from "@noya-app/observable"; import { ExtractSearchParams, Routes } from "./rpc/routes"; import { SerializableRequest, SerializableResponse, TypedBody, UploadProgress, } from "./rpc/types"; export type RPCStatus = "pending" | "resolved" | "rejected" | "streaming"; export type RPCResponseMessage = { id: string } & ( | { type: "end"; response: SerializableResponse } | { type: "chunk"; chunk: string } | { type: "error"; error: any } ); type RequestOptions = Record> = { headers?: Record; body?: string; searchParams?: Partial; onUploadProgress?: (progress: UploadProgress) => void; }; type StreamingHandlers = { onStreamChunk?: (chunk: string) => void; onStreamEnd?: () => Promise; }; // Helper type to extract route key from method and URL type RouteKey = keyof Routes; type RouteMethod = Routes[K]["request"]["options"]["method"]; type RouteURL = Routes[K]["request"]["url"]; export class RPCManager extends Emitter<[SerializableRequest]> { debug: boolean; constructor(options: { isConnected?: boolean; debug?: boolean } = {}) { super(); this.isConnected$ = new Observable(options.isConnected ?? false); this.debug = options.debug ?? false; } isConnected$: Observable; setIsConnected(connected: boolean) { this.isConnected$.set(connected); } requests: Record< string, { request: SerializableRequest; status: RPCStatus; resolve: (response: SerializableResponse) => void; reject: (reason: any) => void; onStreamChunk?: (chunk: string) => void; onStreamEnd?: () => Promise; onUploadProgress?: (progress: UploadProgress) => void; } > = {}; request = async ( request: Routes[K]["request"], options?: RequestOptions & StreamingHandlers & { streaming?: boolean } ): Promise => { if (this.debug) { console.info("[rpcManager] request waiting for connection"); } await this.isConnected$.waitFor(true); if (this.debug) { console.info("[rpcManager] request", request); } const id = uuid(); const requestWithId = { ...request, id }; const { promise, resolve, reject } = Promise.withResolvers(); this.requests[id] = { request: requestWithId, status: "pending", resolve: resolve as (response: SerializableResponse) => void, reject, onStreamChunk: options?.onStreamChunk, onStreamEnd: options?.onStreamEnd, onUploadProgress: options?.onUploadProgress, }; this.emit(requestWithId); return promise; }; requestRoute = async ( route: K, { body, headers, searchParams, onUploadProgress }: RequestOptions> = {} ): Promise => { const [method, baseUrl] = route.split(" ") as [RouteMethod, RouteURL]; // Build URL with search params if provided let url: string = baseUrl; if (searchParams && Object.keys(searchParams).length > 0) { const params = new URLSearchParams(); for (const [key, value] of Object.entries(searchParams)) { if (value !== undefined && value !== null) { params.set(key, String(value)); } } url = `${baseUrl}?${params.toString()}`; } return this.request( { url, options: { method, ...(body ? { body } : {}), ...(headers ? { headers } : {}), }, } as Routes[K]["request"], { onUploadProgress, } ); }; requestStreamingRoute = async ( route: K, { body, headers, onStreamChunk, onStreamEnd, }: RequestOptions & StreamingHandlers ): Promise => { const [method, url] = route.split(" ") as [RouteMethod, RouteURL]; return this.request( { url, options: { method, ...(body ? { body } : {}), ...(headers ? { headers } : {}), }, streaming: true, } as Routes[K]["request"], { onStreamChunk, onStreamEnd, } ); }; handleMessage = (message: RPCResponseMessage) => { const request = this.requests[message.id]; if (this.debug) { console.info("[rpcManager] handleMessage", message); } if (request) { if (message.type === "chunk") { request.status = "streaming"; request.onStreamChunk?.(message.chunk!); } else if (message.type === "end") { request.status = "resolved"; if (request.onStreamEnd) { request.onStreamEnd().then(() => { request.resolve(message.response!); }); } else { request.resolve(message.response!); } } else if (message.type === "error") { request.status = "rejected"; request.reject(message.error); } } else { console.error( "Response without corresponding request", message, this.requests ); } }; getResponseBody = < R extends { body: TypedBody; headers: Record }, >( response: R ): NonNullable => { const headers = Object.fromEntries( Object.entries(response.headers).map(([key, value]) => [ key.toLowerCase(), value, ]) ); const contentType = "content-type" in headers ? headers["content-type"] : null; if (contentType?.includes("application/json")) { return JSON.parse(response.body); } return response.body; }; }