From ec1a50aaf66f218a3e2a10f6eae19559856cb831 Mon Sep 17 00:00:00 2001 From: Daniel Wong Date: Thu, 27 Aug 2026 01:30:23 +0100 Subject: [PATCH 1/2] fix(core): keep the cached text stream observable across a room disconnect --- .changeset/olive-jokes-shave.md | 5 + .../core/src/components/textStream.test.ts | 95 +++++++++++++++++++ packages/core/src/components/textStream.ts | 9 +- 3 files changed, 107 insertions(+), 2 deletions(-) create mode 100644 .changeset/olive-jokes-shave.md create mode 100644 packages/core/src/components/textStream.test.ts diff --git a/.changeset/olive-jokes-shave.md b/.changeset/olive-jokes-shave.md new file mode 100644 index 000000000..09a5618a8 --- /dev/null +++ b/.changeset/olive-jokes-shave.md @@ -0,0 +1,5 @@ +--- +"@livekit/components-core": patch +--- + +fix(core): keep the cached text stream observable across a room disconnect diff --git a/packages/core/src/components/textStream.test.ts b/packages/core/src/components/textStream.test.ts new file mode 100644 index 000000000..7f54f8499 --- /dev/null +++ b/packages/core/src/components/textStream.test.ts @@ -0,0 +1,95 @@ +import { RoomEvent, type Room } from 'livekit-client'; +import { describe, expect, it, vi } from 'vitest'; +import { setupTextStream } from './textStream'; + +/** + * Minimal stand-in for `Room` that mirrors livekit-client's one-handler-per-topic + * contract: registering the same topic twice throws. + */ +function createFakeRoom() { + const handlers = new Map Promise>(); + const disconnectListeners: Array<() => void> = []; + + const registerTextStreamHandler = vi.fn( + (topic: string, handler: (reader: unknown, participantInfo: unknown) => Promise) => { + if (handlers.has(topic)) { + throw new Error(`A text stream handler for topic "${topic}" has already been set.`); + } + handlers.set(topic, handler); + }, + ); + const unregisterTextStreamHandler = vi.fn((topic: string) => { + handlers.delete(topic); + }); + + const room = { + registerTextStreamHandler, + unregisterTextStreamHandler, + on: (event: RoomEvent, listener: () => void) => { + if (event === RoomEvent.Disconnected) disconnectListeners.push(listener); + return room; + }, + } as unknown as Room; + + return { + room, + registerTextStreamHandler, + unregisterTextStreamHandler, + disconnect: () => disconnectListeners.forEach((listener) => listener()), + push: (topic: string, reader: unknown) => handlers.get(topic)?.(reader, { identity: 'agent' }), + }; +} + +function fakeReader(id: string, chunks: string[]) { + return { + info: { id, attributes: {} }, + async *[Symbol.asyncIterator]() { + for (const chunk of chunks) yield chunk; + }, + }; +} + +const TOPIC = 'lk.transcription'; + +describe('setupTextStream', () => { + it('keeps one observable per topic across a disconnect, so a later subscriber shares the handler', () => { + const { room, registerTextStreamHandler, disconnect } = createFakeRoom(); + + // A consumer subscribes while connected. + const first = setupTextStream(room, TOPIC); + first.subscribe().unsubscribe(); // disconnect makes every consumer unsubscribe + + disconnect(); + registerTextStreamHandler.mockClear(); + + // A consumer that mounts *while disconnected* must get the same observable. + // Two observables for one topic both register it on reconnect, and the + // second call throws `A text stream handler ... has already been set`. + const second = setupTextStream(room, TOPIC); + + const subA = first.subscribe(); + const subB = second.subscribe(); + + expect(registerTextStreamHandler).toHaveBeenCalledTimes(1); + expect(second).toBe(first); + + subA.unsubscribe(); + subB.unsubscribe(); + }); + + it('clears buffered streams on disconnect', async () => { + const { room, disconnect, push } = createFakeRoom(); + const emitted: number[] = []; + + const stream = setupTextStream(room, TOPIC); + const sub = stream.subscribe((streams) => emitted.push(streams.length)); + + await push(TOPIC, fakeReader('stream-1', ['hello'])); + await Promise.resolve(); + + disconnect(); + + expect(emitted.at(-1)).toBe(0); + sub.unsubscribe(); + }); +}); diff --git a/packages/core/src/components/textStream.ts b/packages/core/src/components/textStream.ts index 736e6bdcf..cf2226f74 100644 --- a/packages/core/src/components/textStream.ts +++ b/packages/core/src/components/textStream.ts @@ -120,9 +120,14 @@ export function setupTextStream(room: Room, topic: string): Observable { - getObservableCache().delete(cacheKey); textStreams = []; textStreamsSubject.next([]); }); From d70433383c8becb7b35a2ee11cce9bca0a9c16e9 Mon Sep 17 00:00:00 2001 From: Daniel Wong Date: Wed, 2 Sep 2026 14:53:27 +0100 Subject: [PATCH 2/2] fix(core): key the text stream cache on the room instance Address review: instead of keeping a global Map keyed by a generated room id, hold a WeakMap> so entries are collected with the room. The RoomEvent.Disconnected listener is gone entirely - the buffer reset moves into tap({ subscribe }), which share() runs once per subscription window. --- .changeset/olive-jokes-shave.md | 2 +- .../core/src/components/textStream.test.ts | 46 ++++++++----- packages/core/src/components/textStream.ts | 66 +++++-------------- 3 files changed, 49 insertions(+), 65 deletions(-) diff --git a/.changeset/olive-jokes-shave.md b/.changeset/olive-jokes-shave.md index 09a5618a8..c04c0d1f1 100644 --- a/.changeset/olive-jokes-shave.md +++ b/.changeset/olive-jokes-shave.md @@ -2,4 +2,4 @@ "@livekit/components-core": patch --- -fix(core): keep the cached text stream observable across a room disconnect +fix(core): key the text stream observable cache on the room instance so it survives a disconnect diff --git a/packages/core/src/components/textStream.test.ts b/packages/core/src/components/textStream.test.ts index 7f54f8499..27ea8723b 100644 --- a/packages/core/src/components/textStream.test.ts +++ b/packages/core/src/components/textStream.test.ts @@ -34,7 +34,6 @@ function createFakeRoom() { return { room, registerTextStreamHandler, - unregisterTextStreamHandler, disconnect: () => disconnectListeners.forEach((listener) => listener()), push: (topic: string, reader: unknown) => handlers.get(topic)?.(reader, { identity: 'agent' }), }; @@ -49,16 +48,19 @@ function fakeReader(id: string, chunks: string[]) { }; } +/** `from(reader)` walks the async iterator, so emissions land a few microtasks later. */ +const flush = () => new Promise((resolve) => setTimeout(resolve, 0)); + const TOPIC = 'lk.transcription'; describe('setupTextStream', () => { - it('keeps one observable per topic across a disconnect, so a later subscriber shares the handler', () => { + it('reuses one observable per room and topic across disconnect/reconnect', () => { const { room, registerTextStreamHandler, disconnect } = createFakeRoom(); - // A consumer subscribes while connected. + // A consumer subscribes while connected, then the room disconnects and + // `useTextStream` drops the subscription. const first = setupTextStream(room, TOPIC); - first.subscribe().unsubscribe(); // disconnect makes every consumer unsubscribe - + first.subscribe().unsubscribe(); disconnect(); registerTextStreamHandler.mockClear(); @@ -77,19 +79,31 @@ describe('setupTextStream', () => { subB.unsubscribe(); }); - it('clears buffered streams on disconnect', async () => { - const { room, disconnect, push } = createFakeRoom(); - const emitted: number[] = []; - - const stream = setupTextStream(room, TOPIC); - const sub = stream.subscribe((streams) => emitted.push(streams.length)); + it('caches per room instance, so a second room gets its own observable', () => { + const a = createFakeRoom(); + const b = createFakeRoom(); - await push(TOPIC, fakeReader('stream-1', ['hello'])); - await Promise.resolve(); + expect(setupTextStream(a.room, TOPIC)).not.toBe(setupTextStream(b.room, TOPIC)); + }); - disconnect(); + it('starts each subscription window with an empty buffer', async () => { + const { room, push } = createFakeRoom(); - expect(emitted.at(-1)).toBe(0); - sub.unsubscribe(); + const stream = setupTextStream(room, TOPIC); + const before: number[] = []; + const firstSub = stream.subscribe((streams) => before.push(streams.length)); + await push(TOPIC, fakeReader('stream-1', ['hello'])); + await flush(); + expect(before.at(-1)).toBe(1); + firstSub.unsubscribe(); + + // Reconnect: the buffer from the previous window must not leak into this one. + const after: number[] = []; + const secondSub = stream.subscribe((streams) => after.push(streams.length)); + await push(TOPIC, fakeReader('stream-2', ['world'])); + await flush(); + + expect(after).toEqual([1]); + secondSub.unsubscribe(); }); }); diff --git a/packages/core/src/components/textStream.ts b/packages/core/src/components/textStream.ts index cf2226f74..f17227733 100644 --- a/packages/core/src/components/textStream.ts +++ b/packages/core/src/components/textStream.ts @@ -1,4 +1,4 @@ -import { RoomEvent, type Room, type TextStreamInfo } from 'livekit-client'; +import { type Room, type TextStreamInfo } from 'livekit-client'; import { from, scan, Subject, type Observable } from 'rxjs'; import { share, tap } from 'rxjs/operators'; import { ParticipantAgentAttributes } from '../helper'; @@ -9,47 +9,25 @@ export interface TextStreamData { streamInfo: TextStreamInfo; } -// Singleton getters for lazy initialization -let observableCacheInstance: Map> | null = null; -let roomInstanceMapInstance: WeakMap | null = null; -let nextRoomId = 0; +// One observable per room and topic. The outer map is weak so a room that is no +// longer referenced takes its observables with it, while a reused room keeps +// them across connect/disconnect cycles. +const observableCache = new WeakMap>>(); -// Get or create the observable cache -function getObservableCache(): Map> { - if (!observableCacheInstance) { - observableCacheInstance = new Map>(); +function getTopicCache(room: Room): Map> { + let topicCache = observableCache.get(room); + if (!topicCache) { + topicCache = new Map>(); + observableCache.set(room, topicCache); } - return observableCacheInstance; -} - -// Get or create the room instance map -function getRoomInstanceMap(): WeakMap { - if (!roomInstanceMapInstance) { - roomInstanceMapInstance = new WeakMap(); - } - return roomInstanceMapInstance; -} - -// Helper to generate cache key -function getCacheKey(room: Room, topic: string): string { - const instanceMap = getRoomInstanceMap(); - - // Get or create a unique ID for this room instance - let roomId = instanceMap.get(room); - if (!roomId) { - roomId = `room_${nextRoomId++}`; - instanceMap.set(room, roomId); - } - - return `${roomId}:${topic}`; + return topicCache; } export function setupTextStream(room: Room, topic: string): Observable { - const cacheKey = getCacheKey(room, topic); - const observableCache = getObservableCache(); + const topicCache = getTopicCache(room); // Check if we already have an observable for this room and topic - const existingObservable = observableCache.get(cacheKey); + const existingObservable = topicCache.get(topic); if (existingObservable) { return existingObservable; } @@ -63,6 +41,10 @@ export function setupTextStream(room: Room, topic: string): Observable { + // `share()` resets on refcount zero, so this runs once per subscription + // window: on the first subscriber, and again after every reconnect on a + // reused room. Each window starts from an empty buffer. + textStreams = []; room.registerTextStreamHandler(topic, async (reader, participantInfo) => { // Create an observable from the reader const streamObservable = from(reader).pipe( @@ -118,19 +100,7 @@ export function setupTextStream(room: Room, topic: string): Observable { - textStreams = []; - textStreamsSubject.next([]); - }); + topicCache.set(topic, sharedObservable); return sharedObservable; }