From 2da9c856471a6e25222ea45e762dae25eea9e943 Mon Sep 17 00:00:00 2001 From: Neil <4138956+nwparker@users.noreply.github.com> Date: Sat, 30 May 2026 12:40:22 -0700 Subject: [PATCH] Clean up remote runtime subscription listeners --- src/shared/remote-runtime-client.test.ts | 33 +++++++++++++++ src/shared/remote-runtime-client.ts | 52 ++++++++++++++++++++---- 2 files changed, 76 insertions(+), 9 deletions(-) diff --git a/src/shared/remote-runtime-client.test.ts b/src/shared/remote-runtime-client.test.ts index e352fda837b..e708052d807 100644 --- a/src/shared/remote-runtime-client.test.ts +++ b/src/shared/remote-runtime-client.test.ts @@ -57,6 +57,39 @@ describe('subscribeRemoteRuntimeRequest', () => { expect(onError).not.toHaveBeenCalled() subscription.close() }) + + it('detaches subscription socket listeners after close', async () => { + const offSpy = vi.spyOn(WebSocketClient.prototype, 'off') + try { + const server = await createSubscriptionServer() + const onResponse = vi.fn() + const onError = vi.fn() + const onClose = vi.fn() + + const subscription = await subscribeRemoteRuntimeRequest( + server.pairing, + 'terminal.subscribe', + { terminal: 't1' }, + 1000, + { + onResponse, + onError, + onClose + } + ) + + await vi.waitFor(() => expect(onResponse).toHaveBeenCalled()) + subscription.close() + await vi.waitFor(() => expect(onClose).toHaveBeenCalledOnce()) + + const removedEvents = offSpy.mock.calls.map(([event]) => event) + expect(removedEvents).toEqual(expect.arrayContaining(['open', 'error', 'close', 'message'])) + expect(subscription.sendBinary(new Uint8Array([9]))).toBe(false) + expect(onError).not.toHaveBeenCalled() + } finally { + offSpy.mockRestore() + } + }) }) describe('sendRemoteRuntimeRequest', () => { diff --git a/src/shared/remote-runtime-client.ts b/src/shared/remote-runtime-client.ts index a0da7f48ab6..60aa83d3932 100644 --- a/src/shared/remote-runtime-client.ts +++ b/src/shared/remote-runtime-client.ts @@ -341,6 +341,33 @@ export async function subscribeRemoteRuntimeRequest( let settled = false let ws: WebSocket | null = null + const cleanupSocketListeners = (): WebSocket | null => { + const socket = ws + if (!socket) { + return null + } + socket.off('open', onOpen) + socket.off('error', onError) + socket.off('close', onClose) + socket.off('message', onMessage) + ws = null + // Why: startup failures detach Orca callbacks before closing the ws, + // but ws can still emit a late transport error while close is in flight. + if (socket.readyState !== WebSocket.CLOSED) { + socket.on('error', ignoreSettledRemoteRuntimeSocketError) + } + return socket + } + + const closeSocketAfterCleanup = (): void => { + const socket = cleanupSocketListeners() + try { + socket?.close() + } catch { + // ignore best-effort close + } + } + const timeout = setTimeout(() => { fail( new RemoteRuntimeClientError( @@ -379,7 +406,7 @@ export async function subscribeRemoteRuntimeRequest( if (!settled) { settled = true clearTimeout(timeout) - close() + closeSocketAfterCleanup() reject(error) return } @@ -394,27 +421,29 @@ export async function subscribeRemoteRuntimeRequest( return } - ws.once('open', () => { + function onOpen(): void { ws?.send( JSON.stringify({ type: 'e2ee_hello', publicKeyB64: publicKeyToBase64(keyPair.publicKey) }) ) - }) + } - ws.once('error', () => { + function onError(): void { fail( new RemoteRuntimeClientError( 'remote_runtime_unavailable', 'Could not connect to the remote Orca runtime.' ) ) - }) + } - ws.on('close', () => { + function onClose(): void { clearTimeout(timeout) + cleanupSocketListeners() if (!settled) { + settled = true reject( new RemoteRuntimeClientError( 'remote_runtime_unavailable', @@ -424,9 +453,9 @@ export async function subscribeRemoteRuntimeRequest( return } callbacks.onClose?.() - }) + } - ws.on('message', (data, isBinary) => { + function onMessage(data: WebSocket.RawData, isBinary: boolean): void { if (isBinary) { handleBinaryFrame(new Uint8Array(data as Buffer)) return @@ -455,7 +484,12 @@ export async function subscribeRemoteRuntimeRequest( } handleRpcFrame(plaintext) - }) + } + + ws.once('open', onOpen) + ws.once('error', onError) + ws.on('close', onClose) + ws.on('message', onMessage) function handleReadyFrame(frame: string): void { let ready: unknown