Rewrite reader logic, add separate reader for caching MP4 chunks, add more advanced drain iterators

This commit is contained in:
Vanilagy
2024-12-08 20:04:47 +01:00
parent ecfd193ae0
commit 5978110de9
12 changed files with 691 additions and 355 deletions
+109 -14
View File
@@ -24,11 +24,65 @@ export class EncodedVideoChunkDrain {
return this.videoTrack._backing.getNextKeyChunk(chunk);
}
async* chunks(startTimestamp = 0) {
let chunk = await this.getChunk(startTimestamp); // Not necessarily correct if there is no chunk at timestamp 0
while (chunk) {
yield chunk;
chunk = await this.getNextChunk(chunk);
async* chunks(startChunk?: EncodedVideoChunk, endTimestamp = Infinity) {
const chunkQueue: EncodedVideoChunk[] = [];
let { promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers();
let { promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers();
let ended = false;
const timestamps: number[] = [];
// The queue should always be big enough to hold 1 second worth of chunks
const maxQueueSize = () => Math.max(2, timestamps.length);
// The following is the "pump" process that keeps pumping chunks into the queue
void (async () => {
let chunk = startChunk ?? await this.getFirstChunk();
while (chunk && !ended) {
if (chunk.timestamp / 1e6 >= endTimestamp) {
break;
}
if (chunkQueue.length > maxQueueSize()) {
({ promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers());
await queueDequeue;
continue;
}
chunkQueue.push(chunk);
onQueueNotEmpty();
({ promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers());
chunk = await this.getNextChunk(chunk);
}
ended = true;
onQueueNotEmpty();
})();
try {
while (true) {
if (chunkQueue.length > 0) {
yield chunkQueue.shift()!;
const now = performance.now();
timestamps.push(now);
while (timestamps.length > 0 && now - timestamps[0]! >= 1000) {
timestamps.shift();
}
onQueueDequeue();
} else if (!ended) {
await queueNotEmpty;
} else {
break;
}
}
} finally {
ended = true;
onQueueDequeue();
}
}
}
@@ -109,23 +163,31 @@ export class VideoFrameDrain {
return result;
}
async* frames(startTimestamp = 0) {
async* frames(startTimestamp = 0, endTimestamp = Infinity) {
const frameQueue: VideoFrame[] = [];
let firstFrameQueued = false;
let lastFrame: VideoFrame | null = null;
let { promise: queueNonEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers();
let { promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers();
let ended = false;
const decoder = await this.createDecoder((frame) => {
const frameTimestamp = frame.timestamp / 1e6;
if (frameTimestamp >= endTimestamp) {
ended = true;
}
if (ended) {
frame.close();
return;
}
const frameTimestamp = frame.timestamp / 1e6;
if (lastFrame) {
if (frameTimestamp > startTimestamp) {
// We don't know ahead of time what the first frame is. This is because the first frame is the last
// frame whose timestamp is less than or equal to the start timestamp. Therefore we need to wait
// for the first frame after the start timestamp, and then we'll know that the previous frame was
// the first frame.
frameQueue.push(lastFrame);
firstFrameQueued = true;
} else {
@@ -142,11 +204,12 @@ export class VideoFrameDrain {
if (frameQueue.length > 0) {
onQueueNotEmpty();
({ promise: queueNonEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers());
({ promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers());
}
});
const keyChunk = await this.videoTrack._backing.getKeyChunk(startTimestamp);
const keyChunk = await this.videoTrack._backing.getKeyChunk(startTimestamp)
?? await this.videoTrack._backing.getFirstChunk();
if (!keyChunk) {
return;
}
@@ -157,6 +220,28 @@ export class VideoFrameDrain {
void (async () => {
let currentChunk: EncodedVideoChunk | null = keyChunk;
let chunksEndTimestamp = Infinity;
if (endTimestamp < Infinity) {
// When an end timestamp is set, we cannot simply use that for the chunk iterator due to out-of-order
// frames (B-frames). Instead, we'll need to keep decoding chunks until we get a frame that exceeds
// this end time. However, we can still put a bound on it: Since key frames are by definition never
// out of order, we can stop at the first key frame after the end timestamp.
const endFrame = await this.videoTrack._backing.getChunk(endTimestamp);
const endKeyFrame = !endFrame
? null
: endFrame.type === 'key' && endFrame.timestamp / 1e6 === endTimestamp
? endFrame
: await this.videoTrack._backing.getNextKeyChunk(endFrame);
if (endKeyFrame) {
chunksEndTimestamp = endKeyFrame.timestamp / 1e6;
}
}
const chunkDrain = new EncodedVideoChunkDrain(this.videoTrack);
const chunks = chunkDrain.chunks(keyChunk, chunksEndTimestamp);
await chunks.next();
while (currentChunk && !ended) {
decoder.decode(currentChunk);
@@ -164,13 +249,23 @@ export class VideoFrameDrain {
await new Promise(resolve => decoder.addEventListener('dequeue', resolve, { once: true }));
}
const nextChunk = await this.videoTrack._backing.getNextChunk(currentChunk);
currentChunk = nextChunk;
const chunkResult = await chunks.next();
if (chunkResult.done) {
break;
}
currentChunk = chunkResult.value;
}
await chunks.return();
await decoder.flush();
decoder.close();
if (!firstFrameQueued && lastFrame) {
frameQueue.push(lastFrame);
}
decoderIsFlushed = true;
onQueueNotEmpty(); // To unstuck the generator
})();
@@ -180,7 +275,7 @@ export class VideoFrameDrain {
if (frameQueue.length > 0) {
yield frameQueue.shift()!;
} else if (!decoderIsFlushed) {
await queueNonEmpty;
await queueNotEmpty;
} else {
break;
}