Skip to content

Commit 2bbf154

Browse files
fix: close stream controllers on room disconnect (#633)
1 parent af79ced commit 2bbf154

3 files changed

Lines changed: 76 additions & 6 deletions

File tree

.changeset/shy-rivers-punch.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@livekit/rtc-node': patch
3+
---
4+
5+
Close in-progress stream controllers on room disconnect to prevent FD leaks

packages/livekit-rtc/src/room.ts

Lines changed: 31 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -303,16 +303,43 @@ export class Room extends (EventEmitter as new () => TypedEmitter<RoomCallbacks>
303303
return ev.message.case == 'disconnect' && ev.message.value.asyncId == res.asyncId;
304304
});
305305

306+
this.cleanupOnDisconnect();
307+
308+
FfiClient.instance.removeListener(FfiClientEvent.FfiEvent, this.onFfiEvent);
309+
this.removeAllListeners();
310+
}
311+
312+
private cleanupOnDisconnect() {
313+
// Error all in-progress stream controllers to prevent FD leaks.
314+
// Streams that were receiving data but never got a trailer (e.g. the sender
315+
// disconnected mid-transfer) would otherwise keep their ReadableStream open
316+
// indefinitely, leaking the underlying controller and any buffered chunks.
317+
// Using error() instead of close() signals an abnormal termination to consumers.
318+
for (const [, streamController] of this.byteStreamControllers) {
319+
try {
320+
streamController.controller.error(new Error('Disconnected while receiving'));
321+
} catch {
322+
// controller may already be closed or errored
323+
}
324+
}
325+
this.byteStreamControllers.clear();
326+
327+
for (const [, streamController] of this.textStreamControllers) {
328+
try {
329+
streamController.controller.error(new Error('Disconnected while receiving'));
330+
} catch {
331+
// controller may already be closed or errored
332+
}
333+
}
334+
this.textStreamControllers.clear();
335+
306336
// Clear sidPromise before removing listeners so that a reconnect
307337
// doesn't return a stale, permanently-pending promise.
308338
this.sidPromise = undefined;
309339
// Abort all pending FfiClient.waitFor() listeners so they don't leak.
310340
// This causes any in-flight operations (publishData, publishTrack, etc.)
311341
// to reject and clean up their event listeners.
312342
this.disconnectController.abort();
313-
314-
FfiClient.instance.removeListener(FfiClientEvent.FfiEvent, this.onFfiEvent);
315-
this.removeAllListeners();
316343
}
317344

318345
/**
@@ -630,9 +657,7 @@ export class Room extends (EventEmitter as new () => TypedEmitter<RoomCallbacks>
630657
/*} else if (ev.case == 'connected') {
631658
this.emit(RoomEvent.Connected);*/
632659
} else if (ev.case == 'disconnected') {
633-
// Abort pending waitFor() listeners on server-initiated disconnect too,
634-
// not just on explicit disconnect() calls.
635-
this.disconnectController.abort();
660+
this.cleanupOnDisconnect();
636661
this.emit(RoomEvent.Disconnected, ev.value.reason!);
637662
} else if (ev.case == 'reconnecting') {
638663
this.emit(RoomEvent.Reconnecting);

packages/livekit-rtc/src/tests/e2e.test.ts

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -515,6 +515,46 @@ describeE2E('livekit-rtc e2e', () => {
515515
testTimeoutMs * 2,
516516
);
517517

518+
it(
519+
'cleans up stream controllers when disconnecting during an active stream',
520+
async () => {
521+
const { rooms } = await connectTestRooms(2);
522+
const [receivingRoom, sendingRoom] = rooms;
523+
const topic = 'cleanup-stream-topic';
524+
525+
// Register a handler on the receiving side that will intentionally
526+
// NOT fully consume the stream — simulating an abandoned transfer.
527+
let readerReceived = false;
528+
// eslint-disable-next-line @typescript-eslint/no-unused-vars
529+
receivingRoom!.registerTextStreamHandler(topic, async (_reader, _sender) => {
530+
readerReceived = true;
531+
// Deliberately do not call reader.readAll() so the stream stays open
532+
});
533+
534+
// Start sending a text stream but don't close it
535+
const writer = await sendingRoom!.localParticipant!.streamText({ topic });
536+
await writer.write('partial data');
537+
538+
// Wait for the receiving side to get the stream header
539+
await waitFor(() => readerReceived, {
540+
timeoutMs: 5000,
541+
debugName: 'text stream header received',
542+
});
543+
544+
// Disconnect the receiving room while the stream is still open.
545+
// This should close the stream controller without throwing.
546+
await receivingRoom!.disconnect();
547+
548+
// Also close the writer and disconnect the sender
549+
await writer.close();
550+
await sendingRoom!.disconnect();
551+
552+
// If we got here without hanging or throwing, the stream controller
553+
// was properly cleaned up on disconnect.
554+
},
555+
testTimeoutMs,
556+
);
557+
518558
it(
519559
'cleans up track publications when a remote participant disconnects',
520560
async () => {

0 commit comments

Comments
 (0)