diff --git a/mobile/src/transport/rpc-client-runtime-events.test.ts b/mobile/src/transport/rpc-client-runtime-events.test.ts index 04b2a90039e..ce035f63b0b 100644 --- a/mobile/src/transport/rpc-client-runtime-events.test.ts +++ b/mobile/src/transport/rpc-client-runtime-events.test.ts @@ -136,4 +136,30 @@ describe('runtime client-event stream disposal', () => { ]) client.close() }) + + it('releases the replayed registration on the new socket when disposed before its ready', async () => { + const { client, socket: first } = connectReadyClient() + const unsubscribe = client.subscribe('runtime.clientEvents.subscribe', null, () => {}) + const request = sentRequests(first, 'runtime.clientEvents.subscribe')[0]! + emitReady(first, request.id, 'runtime-events:first-socket') + + first.close() + await vi.advanceTimersByTimeAsync(500) + const second = sockets.at(-1)! + expect(second).not.toBe(first) + second.open() + second.receive(JSON.stringify({ type: 'e2ee_ready' })) + second.receive('encrypted:{"type":"e2ee_authenticated"}') + expect(sentRequests(second, 'runtime.clientEvents.subscribe')).toEqual([ + expect.objectContaining({ id: request.id }) + ]) + + unsubscribe() + emitReady(second, request.id, 'runtime-events:second-socket') + + expect(sentRequests(second, 'runtime.clientEvents.unsubscribe')).toEqual([ + expect.objectContaining({ params: { subscriptionId: 'runtime-events:second-socket' } }) + ]) + client.close() + }) }) diff --git a/mobile/src/transport/rpc-client-stream-registry.test.ts b/mobile/src/transport/rpc-client-stream-registry.test.ts index 0e7d561be6e..5e3b7c1ad23 100644 --- a/mobile/src/transport/rpc-client-stream-registry.test.ts +++ b/mobile/src/transport/rpc-client-stream-registry.test.ts @@ -98,4 +98,105 @@ describe('RpcClientStreamRegistry', () => { params: { subscriptionId: 'browser-screencast:page-1:test' } }) }) + + it('releases a replayed browser stream replaced by a new one before its ready', () => { + const { registry, sent } = createRegistry() + registry.subscribe('browser.screencast', { page: 'page-1' }, () => {}) + const replayedId = sent[0]!.id + registry.handleResponse( + streamingResponse(replayedId, { type: 'ready', subscriptionId: 'page-1-old-connection' }) + ) + + registry.markForReplay() + registry.replayAfterAuthentication() + registry.subscribe('browser.screencast', { page: 'page-2' }, () => {}) + const replacementId = sent.at(-1)!.id + registry.handleResponse( + streamingResponse(replayedId, { type: 'ready', subscriptionId: 'page-1-new-connection' }) + ) + registry.handleResponse( + streamingResponse(replacementId, { type: 'ready', subscriptionId: 'page-2' }) + ) + + expect( + sent + .filter((request) => request.method === 'browser.screencast.unsubscribe') + .map((request) => request.params) + ).toEqual([{ subscriptionId: 'page-1-new-connection' }]) + expect(registry.size()).toBe(1) + }) + + describe.each([ + ['runtime.clientEvents.subscribe', 'runtime.clientEvents.unsubscribe', null], + ['browser.screencast', 'browser.screencast.unsubscribe', { page: 'page-1' }] + ])('%s ready id across a replay', (method, unsubscribeMethod, params) => { + function unsubscribes(sent: SentRequest[]): unknown[] { + return sent.filter((request) => request.method === unsubscribeMethod).map((r) => r.params) + } + + function subscribeReady() { + const harness = createRegistry() + const dispose = harness.registry.subscribe(method, params, () => {}) + const requestId = harness.sent[0]!.id + harness.registry.handleResponse( + streamingResponse(requestId, { type: 'ready', subscriptionId: 'old-connection-id' }) + ) + return { ...harness, dispose, requestId } + } + + it('releases the replayed registration when disposed before its new ready', () => { + const { registry, sent, dispose, requestId } = subscribeReady() + + registry.markForReplay() + registry.replayAfterAuthentication() + dispose() + registry.handleResponse( + streamingResponse(requestId, { type: 'ready', subscriptionId: 'new-connection-id' }) + ) + + expect(unsubscribes(sent)).toEqual([{ subscriptionId: 'new-connection-id' }]) + }) + + it('forgets the previous connection id when marked for replay', () => { + const { registry, sent, dispose } = subscribeReady() + + registry.markForReplay() + dispose() + + // A disposal while disconnected has nothing to name on the next connection. + expect(unsubscribes(sent)).toEqual([]) + expect(registry.size()).toBe(0) + }) + + it('still releases a stream cancelled before its first ready', () => { + const { registry, sent } = createRegistry() + const dispose = registry.subscribe(method, params, () => {}) + const requestId = sent[0]!.id + + dispose() + expect(unsubscribes(sent)).toEqual([]) + registry.handleResponse( + streamingResponse(requestId, { type: 'ready', subscriptionId: 'first-id' }) + ) + + expect(unsubscribes(sent)).toEqual([{ subscriptionId: 'first-id' }]) + expect(registry.size()).toBe(0) + }) + + it('sends one unsubscribe however often the stream is disposed', () => { + const { registry, sent, dispose, requestId } = subscribeReady() + + registry.markForReplay() + registry.replayAfterAuthentication() + dispose() + dispose() + registry.handleResponse( + streamingResponse(requestId, { type: 'ready', subscriptionId: 'new-connection-id' }) + ) + dispose() + + expect(unsubscribes(sent)).toEqual([{ subscriptionId: 'new-connection-id' }]) + expect(registry.size()).toBe(0) + }) + }) }) diff --git a/mobile/src/transport/rpc-client-stream-registry.ts b/mobile/src/transport/rpc-client-stream-registry.ts index 68e50a6e6dd..78287be7ee1 100644 --- a/mobile/src/transport/rpc-client-stream-registry.ts +++ b/mobile/src/transport/rpc-client-stream-registry.ts @@ -108,6 +108,8 @@ export class RpcClientStreamRegistry { this.pendingBrowserRequestId = null for (const [id, stream] of this.streams) { stream.sent = false + // The id named a registration on the closed socket; the replay's ready brings the new one. + stream.subscriptionId = undefined this.resetTerminalRouting(id) } }