import { ConnectionEvent, WebSocketConnection, } from "@noya-app/noya-multiplayer-react"; import { AssetStoreRemote } from "../AssetManager"; import { NoyaManager } from "../NoyaManager"; import { ClientToServerMessage } from "../multiplayer"; import { SerializableRequest } from "../rpc/types"; import { fetchRPC, fetchRPCStream } from "./rpcFetch"; import { applyServerMessage, getSyncConfig } from "./syncUtils"; import { SyncURLOption } from "./types"; type EmbeddedRemoteConnectionOptions = { url: SyncURLOption; debug?: boolean; }; export type EmbeddedRemoteConnection = { sendMessage: (message: ClientToServerMessage) => void; destroy: () => void; }; export type EmbeddedRemoteConnectionDependencies = { getSyncConfig?: typeof getSyncConfig; createWebSocketConnection?: ( url: URL, options?: { debug?: boolean; onConnectionEvent?: (event: ConnectionEvent) => void; } ) => WebSocketConnection; fetchRPC?: typeof fetchRPC; fetchRPCStream?: typeof fetchRPCStream; applyServerMessage?: typeof applyServerMessage; }; /** * Minimal remote connection used by the embedded host to proxy child messages * directly to the multiplayer server without routing through the parent's MSM. */ export function createEmbeddedRemoteConnection< S, M extends object, E extends object, MenuT extends string, I extends Record = Record, >( noyaManager: NoyaManager, { url, debug }: EmbeddedRemoteConnectionOptions, dependencies: EmbeddedRemoteConnectionDependencies = {} ): EmbeddedRemoteConnection { const { getSyncConfig: getConfig = getSyncConfig, createWebSocketConnection = ( connectionUrl: URL, options?: { debug?: boolean; onConnectionEvent?: (event: ConnectionEvent) => void; } ) => new WebSocketConnection(connectionUrl, options), fetchRPC: fetchRPCImpl = fetchRPC, fetchRPCStream: fetchRPCStreamImpl = fetchRPCStream, applyServerMessage: applyServerMessageImpl = applyServerMessage, } = dependencies; let wsConnection: WebSocketConnection | undefined; let isOpen = false; const queue: ClientToServerMessage[] = []; const flushQueue = () => { if (!isOpen) return; while (queue.length && wsConnection) { const message = queue.shift(); if (message) { wsConnection.send(message); } } }; const connectionConfig = getConfig(url).then( ({ url: resolvedUrl, tokenString, tokenPayload }) => { const connectionId = tokenPayload.connectionId; const userId = tokenPayload.authId ?? undefined; noyaManager.sharedConnectionDataManager.setCurrentConnectionId( connectionId ); noyaManager.userManager.setCurrentConnectionId(connectionId); noyaManager.userManager.setCurrentUserId(userId); noyaManager.assetManager.assetStore = new AssetStoreRemote( noyaManager.rpcManager ); wsConnection = createWebSocketConnection(resolvedUrl, { debug, onConnectionEvent: (event: ConnectionEvent) => { noyaManager.connectionEventManager.addEvent(event); switch (event.type) { case "stateChange": { isOpen = event.state === "OPEN"; noyaManager.rpcManager.setIsConnected(isOpen); if (isOpen) { flushQueue(); } break; } case "receive": { applyServerMessageImpl(event.message, noyaManager, { pushToMSM: false, }); break; } } }, }); wsConnection.connect(); return { resolvedUrl, tokenString }; } ); const unsubscribeRpc = noyaManager.rpcManager.addListener( async (request: SerializableRequest) => { const { resolvedUrl, tokenString } = await connectionConfig; if (request.streaming) { const stream = fetchRPCStreamImpl(request, { baseUrl: resolvedUrl.toString(), token: tokenString, streaming: true, debug, }); for await (const message of stream) { noyaManager.rpcManager.handleMessage(message); } } else { const response = await fetchRPCImpl(request, { baseUrl: resolvedUrl.toString(), token: tokenString, debug, onUploadProgress: noyaManager.rpcManager.requests[request.id!]?.onUploadProgress, }); noyaManager.rpcManager.handleMessage(response); } } ); const unsubscribeAiInvocation = noyaManager.aiManager.invocationEmitter.addListener(async (invocation) => { const response = await noyaManager.aiManager.callTool(invocation); noyaManager.aiManager.responseEmitter.emit(invocation, response); }); const sendMessage = (message: ClientToServerMessage) => { if (isOpen && wsConnection) { wsConnection.send(message); return; } queue.push(message); }; const destroy = () => { isOpen = false; noyaManager.rpcManager.setIsConnected(false); if (wsConnection) { wsConnection.close(); } else { connectionConfig.then(() => { wsConnection?.close(); }); } unsubscribeRpc(); unsubscribeAiInvocation(); }; return { sendMessage, destroy }; }