Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/audio-mixer-read-timeout.md
Original file line number Diff line number Diff line change
@@ -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.
45 changes: 45 additions & 0 deletions packages/livekit-rtc/src/audio_mixer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<AudioFrame> {
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();
}
});
});
19 changes: 18 additions & 1 deletion packages/livekit-rtc/src/audio_mixer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,9 @@ export class AudioMixer {
private streams: Set<AudioStream>;
private buffers: Map<AudioStream, Int16Array>;
private streamIterators: Map<AudioStream, { next(): Promise<IteratorResult<AudioFrame>> }>;
// 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<AudioStream, Promise<IteratorResult<AudioFrame>>>;
private sampleRate: number;
private numChannels: number;
private chunkSize: number;
Expand All @@ -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 =
Expand Down Expand Up @@ -141,6 +145,7 @@ export class AudioMixer {
this.streams.delete(stream);
this.buffers.delete(stream);
this.streamIterators.delete(stream);
this.pendingReads.delete(stream);
}

/**
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down
Loading