diff --git a/.changeset/audio-mixer-read-timeout.md b/.changeset/audio-mixer-read-timeout.md new file mode 100644 index 00000000..61a9cc31 --- /dev/null +++ b/.changeset/audio-mixer-read-timeout.md @@ -0,0 +1,5 @@ +--- +'@livekit/rtc-node': patch +--- + +Fix `AudioMixer` dropping audio from a stream whose read times out. The mixer now keeps the pending read and awaits it again, instead of issuing a new `next()` and discarding the frame the first read later resolves with. diff --git a/packages/livekit-rtc/src/audio_mixer.test.ts b/packages/livekit-rtc/src/audio_mixer.test.ts index d4dfbfaa..f852a22e 100644 --- a/packages/livekit-rtc/src/audio_mixer.test.ts +++ b/packages/livekit-rtc/src/audio_mixer.test.ts @@ -266,4 +266,49 @@ describe('AudioMixer', () => { console.warn = originalWarn; } }); + + it('does not drop a frame that arrives after a read timeout', async () => { + const sampleRate = 48000; + const numChannels = 1; + const samplesPerChannel = 480; + const mixer = new AudioMixer(sampleRate, numChannels, { + blocksize: samplesPerChannel, + streamTimeoutMs: 20, + }); + + const makeFrame = (value: number) => + new AudioFrame( + new Int16Array(numChannels * samplesPerChannel).fill(value), + sampleRate, + numChannels, + samplesPerChannel, + ); + + // The first frame arrives well after the mixer's read timeout. + async function* lateStream(): AsyncGenerator { + await new Promise((resolve) => setTimeout(resolve, 100)); + yield makeFrame(111); + yield makeFrame(222); + } + + const originalWarn = console.warn; + console.warn = () => {}; + + try { + mixer.addStream(lateStream()); + + const values: number[] = []; + for await (const frame of mixer) { + const value = frame.data[0]!; + if (value !== 0) { + values.push(value); + } + } + + expect(values).toEqual([111, 222]); + } finally { + console.warn = originalWarn; + await mixer.aclose(); + } + }); }); diff --git a/packages/livekit-rtc/src/audio_mixer.ts b/packages/livekit-rtc/src/audio_mixer.ts index e2aaef84..270849cb 100644 --- a/packages/livekit-rtc/src/audio_mixer.ts +++ b/packages/livekit-rtc/src/audio_mixer.ts @@ -70,6 +70,9 @@ export class AudioMixer { private streams: Set; private buffers: Map; private streamIterators: Map> }>; + // A read that timed out is still in flight. It is kept here and awaited again on the next + // pass, so the frame it eventually resolves with is not lost. + private pendingReads: Map>>; private sampleRate: number; private numChannels: number; private chunkSize: number; @@ -91,6 +94,7 @@ export class AudioMixer { this.streams = new Set(); this.buffers = new Map(); this.streamIterators = new Map(); + this.pendingReads = new Map(); this.sampleRate = sampleRate; this.numChannels = numChannels; this.chunkSize = @@ -141,6 +145,7 @@ export class AudioMixer { this.streams.delete(stream); this.buffers.delete(stream); this.streamIterators.delete(stream); + this.pendingReads.delete(stream); } /** @@ -311,13 +316,24 @@ export class AudioMixer { // Accumulate data until we have at least chunkSize samples while (buf.length < this.chunkSize * this.numChannels && !exhausted && !this.closed) { try { - const result = await this.timeoutRace(iterator.next(), this.streamTimeoutMs); + let read = this.pendingReads.get(stream); + if (!read) { + read = iterator.next(); + this.pendingReads.set(stream, read); + } + const result = await this.timeoutRace(read, this.streamTimeoutMs); if (result === 'timeout') { + // Keep the pending read: issuing another next() would leave this one to resolve + // unobserved and its frame would be dropped. console.warn(`AudioMixer: stream timeout after ${this.streamTimeoutMs}ms`); break; } + if (this.pendingReads.get(stream) === read) { + this.pendingReads.delete(stream); + } + if (result.done) { exhausted = true; break; @@ -339,6 +355,7 @@ export class AudioMixer { buf = combined; } } catch (error) { + this.pendingReads.delete(stream); console.error(`AudioMixer: Error reading from stream:`, error); exhausted = true; break;