import { ConnectionEvent, WebSocketConnection, } from "@noya-app/noya-multiplayer-react"; import { AssetStoreRemote } from "../AssetManager"; import { ClientToServerMessage } from "../multiplayer"; import { SerializableRequest } from "../rpc/types"; import { fetchRPC, fetchRPCStream } from "./rpcFetch"; import { applyServerMessage, getOptionsObject, getSyncConfig, } from "./syncUtils"; import { SyncAdapter, WebSocketSyncOptions } from "./types"; type CreateWebSocketConnection = ( url: URL, options?: { debug?: boolean; onConnectionEvent?: (event: ConnectionEvent) => void; } ) => WebSocketConnection; type WebSocketSyncDependencies = { createWebSocketConnection?: CreateWebSocketConnection; }; let count = 0; export function webSocketSync< S, M extends object, E extends object, MenuT extends string = string, I extends Record = Record, >( options: WebSocketSyncOptions, dependencies: WebSocketSyncDependencies = {} ): SyncAdapter { const { debug = false } = getOptionsObject(options); const createConnection: CreateWebSocketConnection = dependencies.createWebSocketConnection ?? ((connectionUrl, connectionOptions) => new WebSocketConnection(connectionUrl, connectionOptions)); return ({ noyaManager }) => { const syncId = count++; let connectionId: string | undefined; let userId: string | undefined; let wsConnection: WebSocketConnection | undefined; const connectionConfig = getSyncConfig(options).then( ({ url, tokenString, tokenPayload }) => { connectionId = tokenPayload.connectionId; userId = tokenPayload.authId ?? undefined; noyaManager.sharedConnectionDataManager.setCurrentConnectionId( connectionId ); noyaManager.userManager.setCurrentConnectionId(connectionId); noyaManager.userManager.setCurrentUserId(userId); noyaManager.multiplayerStateManager.setPolicyAuthContext({ id: userId ?? null, access: tokenPayload.access, }); noyaManager.assetManager.assetStore = new AssetStoreRemote( noyaManager.rpcManager ); wsConnection = createConnection(url, { debug, onConnectionEvent: (event) => { if (debug) { console.info(`[webSocketSync] Connection event ${syncId}`, event); } switch (event.type) { case "stateChange": { if (event.state === "OPEN") { noyaManager.multiplayerStateManager.connect(); noyaManager.rpcManager.setIsConnected(true); } else if (event.state === "CLOSED") { noyaManager.multiplayerStateManager.disconnect(); noyaManager.rpcManager.setIsConnected(false); } break; } case "receive": { noyaManager.connectionEventManager.addEvent(event); const message = event.message; // Apply to adapters + optionally push to MSM (full mode only) applyServerMessage(message, noyaManager, { pushToMSM: true, }); break; } case "send": { noyaManager.connectionEventManager.addEvent(event); break; } case "error": { noyaManager.connectionEventManager.addEvent(event); break; } } }, }); wsConnection.connect(); return { url, tokenString, tokenPayload, wsConnection }; } ); const unsubscribeMultiplayerMessage = noyaManager.multiplayerStateManager.messageEmitter.addListener( (message: ClientToServerMessage) => { connectionConfig.then(({ wsConnection }) => { wsConnection?.send(message); }); } ); const unsubscribeSharedConnectionData = noyaManager.sharedConnectionDataManager.currentConnectionDataEmitter.addListener( (sharedConnectionData: E | undefined) => { if (!connectionId) return; const message: ClientToServerMessage = { type: "setSharedConnectionData", connectionId, data: sharedConnectionData, }; connectionConfig.then(({ wsConnection }) => { wsConnection?.send(message); }); } ); const unsubscribeRpc = noyaManager.rpcManager.addListener( async (request: SerializableRequest) => { const { url: baseUrl, tokenString } = await connectionConfig; if (request.streaming) { const stream = fetchRPCStream(request, { baseUrl: baseUrl.toString(), token: tokenString, streaming: true, debug, }); for await (const message of stream) { noyaManager.rpcManager.handleMessage(message); } } else { const response = await fetchRPC(request, { baseUrl: baseUrl.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); } ); if (debug) { console.info(`[webSocketSync] Created ${syncId}`); } return () => { if (debug) { console.info(`[webSocketSync] Destroying ${syncId}`); } // Try to close the connection even without waiting for the connection config // so that we can close as soon as possible if (wsConnection) { if (debug) { console.info(`[webSocketSync] Closing connection ${syncId}`); } wsConnection.close(); } else { connectionConfig.then(({ wsConnection }) => { if (debug) { console.info(`[webSocketSync] Closing connection ${syncId}`); } wsConnection?.close(); }); } unsubscribeMultiplayerMessage(); unsubscribeSharedConnectionData(); unsubscribeRpc(); unsubscribeAiInvocation(); }; }; }