From 88e322c8d9dd0da9ce7131656930d58d187c524b Mon Sep 17 00:00:00 2001 From: lukasIO Date: Wed, 2 Sep 2026 12:32:57 +0200 Subject: [PATCH 1/2] fix: wait for ReconnectResponse to arrive before declaring signal reconnected --- src/api/SignalClient.test.ts | 24 ++++---- src/api/SignalClient.ts | 107 ++++++++++++++++++----------------- 2 files changed, 68 insertions(+), 63 deletions(-) diff --git a/src/api/SignalClient.test.ts b/src/api/SignalClient.test.ts index 8a4c71ee09..b4c9bc1dda 100644 --- a/src/api/SignalClient.test.ts +++ b/src/api/SignalClient.test.ts @@ -236,7 +236,7 @@ describe('SignalClient.connect', () => { expect(signalClient.currentState).toBe(SignalConnectionState.CONNECTED); }); - it('should handle reconnect with non-reconnect message (edge case)', async () => { + it('drops stray messages during reconnect and waits for the reconnect response', async () => { // First, initial connection const joinResponse = createJoinResponse(); const joinSignalResponse = createSignalResponse('join', joinResponse); @@ -247,19 +247,24 @@ describe('SignalClient.connect', () => { await signalClient.join('wss://test.livekit.io', 'test-token', defaultOptions); - // Setup reconnect with non-reconnect message (e.g., participant update) + // Server sends a stray update first, then the actual reconnect response. + // The stray should be dropped with a warning; the reconnect response should resolve. const updateSignalResponse = createSignalResponse('update', { participants: [] }); - const reconnectMockReadable = createMockReadableStream([updateSignalResponse]); + const reconnectResponse = new ReconnectResponse({ iceServers: [] }); + const reconnectSignalResponse = createSignalResponse('reconnect', reconnectResponse); + const reconnectMockReadable = createMockReadableStream([ + updateSignalResponse, + reconnectSignalResponse, + ]); const reconnectMockConnection = createMockConnection(reconnectMockReadable); mockWebSocketStream({ connection: reconnectMockConnection }); const result = await signalClient.reconnect('wss://test.livekit.io', 'test-token', 'sid-123'); - // This is an edge case: reconnect resolves with undefined when non-reconnect message is received - expect(result).toBeUndefined(); + expect(result).toEqual(reconnectResponse); expect(signalClient.currentState).toBe(SignalConnectionState.CONNECTED); - }, 1000); + }); }); describe('Failure Case - Timeout', () => { @@ -1026,7 +1031,7 @@ describe('SignalClient.validateFirstMessage', () => { } }); - it('should accept non-reconnect message during reconnecting state', async () => { + it('should reject non-reconnect message during reconnecting state', async () => { // First establish a connection const joinResponse = createJoinResponse(); const joinSignalResponse = createSignalResponse('join', joinResponse); @@ -1044,9 +1049,8 @@ describe('SignalClient.validateFirstMessage', () => { const validateMethod = (signalClient as any).validateFirstMessage; if (validateMethod) { const result = validateMethod.call(signalClient, updateSignalResponse, true); - expect(result.isValid).toBe(true); - expect(result.response).toBeUndefined(); - expect(result.shouldProcessFirstMessage).toBe(true); + expect(result.isValid).toBe(false); + expect(result.error).toBeInstanceOf(ConnectionError); } }); diff --git a/src/api/SignalClient.ts b/src/api/SignalClient.ts index 72f58d684c..7207e8584a 100644 --- a/src/api/SignalClient.ts +++ b/src/api/SignalClient.ts @@ -581,38 +581,61 @@ export class SignalClient { this.streamWriter = connection.writable.getWriter(); // wsTimeout only guarded the upgrade; guard the first-message read with - // its own timeout so a silent server can't hang join() forever. - let firstMessage: ReadableStreamReadResult; - let firstMessageTimeout: ReturnType | undefined; + // its own deadline so a silent server can't hang join() forever. + // During reconnect, drop any stray messages that arrive before the + // ReconnectResponse so we do not declare the channel reconnected on the + // wrong signal. + let firstSignalResponse: SignalResponse; try { - firstMessage = await Promise.race([ - signalReader.read(), - new Promise((_, rejectRead) => { - firstMessageTimeout = setTimeout(() => { - rejectRead( - ConnectionError.timeout( - 'signal connection timed out while waiting for the first message', - ), - ); - }, JOIN_RESPONSE_TIMEOUT); - }), - ]); + const deadline = Date.now() + JOIN_RESPONSE_TIMEOUT; + // eslint-disable-next-line no-constant-condition + while (true) { + const remaining = deadline - Date.now(); + if (remaining <= 0) { + throw ConnectionError.timeout( + 'signal connection timed out while waiting for the first message', + ); + } + let readTimeout: ReturnType | undefined; + const readResult = await Promise.race([ + signalReader.read(), + new Promise((_, rejectRead) => { + readTimeout = setTimeout(() => { + rejectRead( + ConnectionError.timeout( + 'signal connection timed out while waiting for the first message', + ), + ); + }, remaining); + }), + ]); + clearTimeout(readTimeout); + if (!readResult.value) { + throw ConnectionError.internal('no message received as first message'); + } + const parsed = parseSignalResponse(readResult.value); + if ( + this.lifecycleState === 'reconnecting' && + parsed.message?.case !== 'reconnect' && + parsed.message?.case !== 'leave' + ) { + this.log.warn('dropping signal message while awaiting reconnect response', { + messageCase: parsed.message?.case, + }); + continue; + } + firstSignalResponse = parsed; + break; + } } catch (e) { - // No first message in time: release the reader and tear down the ws + // No usable first message in time: release the reader and tear down the ws // so we surface the timeout instead of leaking an open connection. signalReader.releaseLock(); reject(e); this.close(); return; - } finally { - clearTimeout(firstMessageTimeout); } signalReader.releaseLock(); - if (!firstMessage.value) { - throw ConnectionError.internal('no message received as first message'); - } - - const firstSignalResponse = parseSignalResponse(firstMessage.value); // Validate the first message const validation = this.validateFirstMessage( @@ -643,11 +666,7 @@ export class SignalClient { } } - // Handle successful connection - const firstMessageToProcess = validation.shouldProcessFirstMessage - ? firstSignalResponse - : undefined; - this.handleSignalConnected(connection, wsTimeout, attemptId, firstMessageToProcess); + this.handleSignalConnected(connection, wsTimeout, attemptId); resolve(validation.response); } catch (e) { reject(e); @@ -660,13 +679,7 @@ export class SignalClient { }); } - async startReadingLoop( - signalReader: ReadableStreamDefaultReader, - firstMessage?: SignalResponse, - ) { - if (firstMessage) { - this.handleSignalResponse(firstMessage); - } + async startReadingLoop(signalReader: ReadableStreamDefaultReader) { const attemptId = this.attemptId; while (true) { if (this.signalLatency) { @@ -1184,7 +1197,6 @@ export class SignalClient { connection: WebSocketConnection, timeoutHandle: ReturnType, attemptId: number, - firstMessage?: SignalResponse, ) { clearTimeout(timeoutHandle); const established = this.sendLifecycleInput( @@ -1205,7 +1217,7 @@ export class SignalClient { } this.log.info('signal connected'); this.startPingInterval(); - this.startReadingLoop(connection.readable.getReader(), firstMessage); + this.startReadingLoop(connection.readable.getReader()); } /** @@ -1222,7 +1234,6 @@ export class SignalClient { isValid: boolean; response?: JoinResponse | ReconnectResponse; error?: ConnectionError; - shouldProcessFirstMessage?: boolean; } { if (firstSignalResponse.message?.case === 'join') { return { @@ -1231,22 +1242,12 @@ export class SignalClient { }; } else if ( this.lifecycleState === 'reconnecting' && - firstSignalResponse.message?.case !== 'leave' + firstSignalResponse.message?.case === 'reconnect' ) { - if (firstSignalResponse.message?.case === 'reconnect') { - return { - isValid: true, - response: firstSignalResponse.message.value, - }; - } else { - // in reconnecting, any message received means signal reconnected and we still need to process it - this.log.debug('declaring signal reconnected without reconnect response received'); - return { - isValid: true, - response: undefined, - shouldProcessFirstMessage: true, - }; - } + return { + isValid: true, + response: firstSignalResponse.message.value, + }; } else if (this.isEstablishingConnection && firstSignalResponse.message?.case === 'leave') { return { isValid: false, From f0387f65e7e24818e727a8969b1079a662b58407 Mon Sep 17 00:00:00 2001 From: lukasIO Date: Wed, 2 Sep 2026 13:41:45 +0200 Subject: [PATCH 2/2] clarify --- src/api/SignalClient.e2e.test.ts | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/api/SignalClient.e2e.test.ts b/src/api/SignalClient.e2e.test.ts index d474100a32..702e8232ba 100644 --- a/src/api/SignalClient.e2e.test.ts +++ b/src/api/SignalClient.e2e.test.ts @@ -174,9 +174,7 @@ describe.skipIf(!!unavailable)('SignalClient e2e', () => { it('rejects a reconnect when a leave arrives as the first message', async () => { // Reconnecting into a room that sends leave-first exercises the client's - // first-message validation while in RECONNECTING (path-independent: it - // doesn't rely on the mock detecting reconnect, which v1 hides inside the - // gzipped join_request the mock ignores). + // first-message validation while in RECONNECTING. await join('happy'); const token = await createToken({ signal: 'leave_first_message' }); const err = await client.reconnect(serverUrl, token, 'RM_session').then(