import { describe, expect, it, jest } from "bun:test"; import { NoyaManager } from "../NoyaManager"; import { ClientToServerMessage, ServerToClientMessage } from "../multiplayer"; import { createEmbeddedRemoteConnection, EmbeddedRemoteConnectionDependencies, } from "../sync/embeddedRemoteConnection"; type FetchRPCFn = NonNullable< EmbeddedRemoteConnectionDependencies["fetchRPC"] >; type FetchRPCStreamFn = NonNullable< EmbeddedRemoteConnectionDependencies["fetchRPCStream"] >; type ApplyServerMessageFn = NonNullable< EmbeddedRemoteConnectionDependencies["applyServerMessage"] >; type SetupOptions = { dependencyOverrides?: Partial>; }; function createTestManager() { return new NoyaManager(null); } const flushAsync = () => new Promise((resolve) => setTimeout(resolve, 0)); function createFetchRPCStub(responseId: string) { const mock = jest.fn( async ( _request?: Parameters[0], _options?: Parameters[1] ) => ({ id: responseId, type: "end" as const, response: { status: 200, statusText: "OK", headers: {}, body: "", }, }) ); const fn = (( request: Parameters[0], options: Parameters[1] ) => mock(request, options)) as unknown as FetchRPCFn; return { fn, mock }; } function createFetchRPCStreamStub(responseId: string) { const mock = jest.fn(async function* ( _request?: Parameters[0], _options?: Parameters[1] ) { yield { id: responseId, type: "chunk" as const, chunk: "chunk-data" }; yield { id: responseId, type: "end" as const, response: { status: 200, statusText: "OK", headers: {}, body: "", }, }; }); const fn = (( request: Parameters[0], options: Parameters[1] ) => mock(request, options)) as unknown as FetchRPCStreamFn; return { fn, mock }; } function createApplyServerMessageStub() { const mock = jest.fn( ( _message: ServerToClientMessage, _noyaManager: NoyaManager, _options?: { pushToMSM?: boolean } ) => {} ); const fn = (( message: Parameters[0], manager: Parameters[1], options: Parameters[2] ) => mock(message, manager, options)) as unknown as ApplyServerMessageFn; return { fn, mock }; } function setup(options: SetupOptions = {}) { const manager = createTestManager(); const connect = jest.fn(); const send = jest.fn(); const close = jest.fn(); let onConnectionEvent: ((event: any) => void) | undefined; const { fn: fetchRPC, mock: fetchRPCMock } = createFetchRPCStub("req"); const { fn: fetchRPCStream, mock: fetchRPCStreamMock } = createFetchRPCStreamStub("req"); const { fn: applyServerMessage, mock: applyServerMessageMock } = createApplyServerMessageStub(); const defaultDependencies: EmbeddedRemoteConnectionDependencies = { getSyncConfig: async () => ({ url: new URL("https://example.com/multiplayer?token=test-token"), tokenString: "test-token", tokenPayload: { connectionId: "connection-1", authId: "user-1", fileId: "file-1", access: "write", baseUrl: "https://api.example.com", }, }), createWebSocketConnection: (_url, opts) => { onConnectionEvent = opts?.onConnectionEvent; return { connect, send, close, } as any; }, fetchRPC, fetchRPCStream, applyServerMessage, }; const dependencies = { ...defaultDependencies, ...options.dependencyOverrides, }; const remote = createEmbeddedRemoteConnection( manager, { url: "https://example.com/multiplayer?token=test-token" }, dependencies ); const emitConnectionEvent = (event: any) => { onConnectionEvent?.(event); }; return { manager, remote, dependencies, mocks: { fetchRPC: fetchRPCMock, fetchRPCStream: fetchRPCStreamMock, applyServerMessage: applyServerMessageMock, }, connect, send, close, emitConnectionEvent, }; } describe("createEmbeddedRemoteConnection", () => { it("queues messages until the websocket opens", async () => { const { remote, send, emitConnectionEvent } = setup(); const payload: ClientToServerMessage = { type: "ping", id: "p1" }; remote.sendMessage(payload); expect(send).not.toHaveBeenCalled(); await flushAsync(); emitConnectionEvent({ type: "stateChange", state: "OPEN" }); expect(send).toHaveBeenCalledWith(payload); }); it("applies server messages without pushing to MSM", async () => { const { mocks, emitConnectionEvent, manager } = setup(); await flushAsync(); const serverMessage: ServerToClientMessage = { type: "pong", id: "123", }; emitConnectionEvent({ type: "receive", message: serverMessage }); expect(mocks.applyServerMessage).toHaveBeenCalledWith( serverMessage, manager, { pushToMSM: false } ); }); it("bridges RPC requests via fetchRPC", async () => { const { fn: fetchRPC, mock: fetchRPCMock } = createFetchRPCStub("rpc-1"); const { manager, emitConnectionEvent } = setup({ dependencyOverrides: { fetchRPC }, }); const originalHandleMessage = manager.rpcManager.handleMessage; const handleMessageSpy = jest.fn<(message: any) => void>(); manager.rpcManager.handleMessage = handleMessageSpy as any; await flushAsync(); emitConnectionEvent({ type: "stateChange", state: "OPEN" }); const request = { id: "rpc-1", url: "/rpc", options: { method: "GET" }, }; manager.rpcManager.emit(request); await flushAsync(); const matchingCall = fetchRPCMock.mock.calls.find( ([calledRequest]) => calledRequest === request ) as [typeof request, { token?: string; baseUrl?: string }] | undefined; expect(matchingCall).toBeTruthy(); const [, options] = matchingCall!; expect(options?.token).toBe("test-token"); expect(options?.baseUrl).toBe( "https://example.com/multiplayer?token=test-token" ); expect( handleMessageSpy.mock.calls.some( ([message]) => message?.id === "rpc-1" && message?.type === "end" ) ).toBe(true); manager.rpcManager.handleMessage = originalHandleMessage; }); it("bridges streaming RPC requests via fetchRPCStream", async () => { const chunks: any[] = [ { id: "rpc-stream", type: "chunk", chunk: "chunk-data" }, { id: "rpc-stream", type: "end", response: { status: 200, statusText: "OK", headers: {}, body: "", }, }, ]; type FetchRPCStreamFn = NonNullable< EmbeddedRemoteConnectionDependencies["fetchRPCStream"] >; const fetchRPCStreamMock = jest.fn(async function* ( _request?: Parameters[0], _options?: Parameters[1] ) { for (const chunk of chunks) { yield chunk; } }); const fetchRPCStream = (( request: Parameters[0], options: Parameters[1] ) => fetchRPCStreamMock(request, options)) as unknown as FetchRPCStreamFn; const { manager, emitConnectionEvent } = setup({ dependencyOverrides: { fetchRPCStream }, }); const originalHandleMessage = manager.rpcManager.handleMessage; const handleMessageSpy = jest.fn<(message: any) => void>(); manager.rpcManager.handleMessage = handleMessageSpy as any; await flushAsync(); emitConnectionEvent({ type: "stateChange", state: "OPEN" }); const request = { id: "rpc-stream", url: "/rpc", streaming: true, options: { method: "GET" }, }; manager.rpcManager.emit(request); await flushAsync(); expect(fetchRPCStreamMock).toHaveBeenCalled(); const matchingCalls = handleMessageSpy.mock.calls.filter( ([message]) => message?.id === "rpc-stream" ); expect(matchingCalls).toHaveLength(chunks.length); chunks.forEach((chunk) => { expect( matchingCalls.some( ([message]) => message?.id === chunk.id && message?.type === chunk.type && ("chunk" in chunk ? message?.chunk === chunk.chunk : true) ) ).toBe(true); }); manager.rpcManager.handleMessage = originalHandleMessage; }); it("closes the websocket on destroy", async () => { const { remote, close } = setup(); await flushAsync(); remote.destroy(); expect(close).toHaveBeenCalled(); }); });