Make cursor initialization instant, fix some code

This commit is contained in:
Vanilagy
2025-12-15 01:31:35 +01:00
parent fbcb4dcb95
commit 9222d943ad
2 changed files with 136 additions and 77 deletions
+76 -40
View File
@@ -360,7 +360,7 @@ export abstract class SampleCursor<
transform: SampleTransformer<Sample, TransformedSample>;
autoClose: boolean;
pumpRunning = false;
decoder: DecoderWrapper<Sample>;
decoder: DecoderWrapper<Sample> | null = null;
currentRaw: Sample | null = null;
current: TransformedSample | null = null;
sampleQueue: Sample[] = [];
@@ -389,6 +389,7 @@ export abstract class SampleCursor<
pumpsStarted: 0,
seekPackets: [] as (EncodedPacket | null)[],
decodedPackets: [] as EncodedPacket[],
throwInDecoderInit: false,
throwInPump: false,
throwDecoderError: false,
};
@@ -397,19 +398,29 @@ export abstract class SampleCursor<
return this._closed;
}
protected constructor(
abstract initDecoder(): Promise<DecoderWrapper<Sample>>;
constructor(
reader: PacketReader,
decoder: DecoderWrapper<Sample>,
options: SampleCursorOptions<Sample, TransformedSample>,
) {
// todo options validation
this.packetReader = reader;
this.decoder = decoder;
this.packetCursor = new PacketCursor(reader);
this.options = options;
this.autoClose = options.autoClose ?? true;
this.transform = options.transform ?? (sample => sample as unknown as TransformedSample);
reader.track.input._openSampleCursors.add(this);
const lock = this.pumpMutex.lock();
assert(!lock.pending);
void this.initDecoder()
.then(decoder => this.decoder = decoder)
.catch(error => this.closeWithError(error, false))
.finally(() => lock.release());
}
async onDecoderSample(sample: Sample) {
@@ -521,13 +532,13 @@ export abstract class SampleCursor<
if (lock.pending) await lock.ready;
}
this._ensureNotClosed();
using deferred = defer(() => {
lock?.release();
this.predictedRequests--;
});
this._ensureNotClosed();
const targetPacket = targetPacketPromise instanceof Promise
? await targetPacketPromise
: targetPacketPromise;
@@ -842,22 +853,28 @@ export abstract class SampleCursor<
closePromise: Promise<void> | null = null;
close() {
return this.closePromise ??= (async () => {
this.predictedRequests++;
this.packetReader.track.input._openSampleCursors.delete(this);
return this.closePromise ??= this.closeInternal();
}
using lock = this.pumpMutex.lock();
async closeInternal(doLock = true) {
this.predictedRequests++;
this.packetReader.track.input._openSampleCursors.delete(this);
let lock: AsyncMutexLock | null = null;
if (doLock) {
lock = this.pumpMutex.lock();
if (lock.pending) await lock.ready;
}
using _ = defer(() => lock?.release());
this._closed = true;
this._closed = true;
if (this.pumpRunning) {
await this.stopPump();
}
if (this.pumpRunning) {
await this.stopPump();
}
this.setCurrent(null, null);
this.decoder.close();
})();
this.setCurrent(null, null);
this.decoder?.close();
}
[Symbol.asyncDispose]() {
@@ -874,10 +891,11 @@ export abstract class SampleCursor<
}
async runPump() {
try {
assert(this.packetCursor.current);
assert(this.pumpTarget);
assert(this.packetCursor.current);
assert(this.pumpTarget);
assert(this.decoder);
try {
this.pumpRunning = true;
if (this.debugInfo.enabled) {
@@ -960,7 +978,7 @@ export abstract class SampleCursor<
}
}
closeWithError(error: unknown) {
closeWithError(error: unknown, doLock?: boolean) {
if (this._closed) {
return;
}
@@ -977,7 +995,7 @@ export abstract class SampleCursor<
};
this.pendingRequests.forEach(uh);
return this.close();
return this.closeInternal(doLock);
}
closeWithErrorAndThrow(error: unknown): never {
@@ -987,15 +1005,21 @@ export abstract class SampleCursor<
}
export class VideoSampleCursor<TransformedSample = VideoSample> extends SampleCursor<VideoSample, TransformedSample> {
static async init<T = VideoSample>(
override packetReader!: PacketReader<InputVideoTrack>;
constructor(
reader: PacketReader<InputVideoTrack>,
options: SampleCursorOptions<VideoSample, T> = {},
options: SampleCursorOptions<VideoSample, TransformedSample> = {},
) {
if (!(reader instanceof PacketReader) || !(reader.track instanceof InputVideoTrack)) {
throw new TypeError('reader must be a PacketReader for an InputVideoTrack.');
}
const track = reader.track;
super(reader, options);
}
override async initDecoder(): Promise<DecoderWrapper<VideoSample>> {
const track = this.packetReader.track;
if (!(await track.canDecode())) {
throw new Error(
@@ -1004,36 +1028,45 @@ export class VideoSampleCursor<TransformedSample = VideoSample> extends SampleCu
);
}
if (this.debugInfo.enabled && this.debugInfo.throwInDecoderInit) {
throw new Error('Fake decoder init error!');
}
const decoderConfig = await track.getDecoderConfig();
assert(decoderConfig);
assert(track.codec);
const decoder = new VideoDecoderWrapper(
sample => cursor.onDecoderSample(sample),
error => cursor.onDecoderError(error),
sample => this.onDecoderSample(sample),
error => this.onDecoderError(error),
track.codec,
decoderConfig,
track.rotation,
track.timeResolution,
);
decoder.onDequeue = () => cursor.onDecoderDequeue();
decoder.onDequeue = () => this.onDecoderDequeue();
const cursor: VideoSampleCursor<T> = new VideoSampleCursor(reader, decoder, options);
return cursor;
return decoder;
}
}
export class AudioSampleCursor<TransformedSample = AudioSample> extends SampleCursor<AudioSample, TransformedSample> {
static async init<T = AudioSample>(
override packetReader!: PacketReader<InputAudioTrack>;
constructor(
reader: PacketReader<InputAudioTrack>,
options: SampleCursorOptions<AudioSample, T> = {},
options: SampleCursorOptions<AudioSample, TransformedSample> = {},
) {
if (!(reader instanceof PacketReader) || !(reader.track instanceof InputAudioTrack)) {
throw new TypeError('reader must be a PacketReader for an InputAudioTrack.');
}
const track = reader.track;
super(reader, options);
}
override async initDecoder(): Promise<DecoderWrapper<AudioSample>> {
const track = this.packetReader.track;
if (!(await track.canDecode())) {
throw new Error(
@@ -1042,6 +1075,10 @@ export class AudioSampleCursor<TransformedSample = AudioSample> extends SampleCu
);
}
if (this.debugInfo.enabled && this.debugInfo.throwInDecoderInit) {
throw new Error('Fake decoder init error!');
}
const codec = track.codec;
const decoderConfig = await track.getDecoderConfig();
assert(codec && decoderConfig);
@@ -1049,23 +1086,22 @@ export class AudioSampleCursor<TransformedSample = AudioSample> extends SampleCu
let decoder: AudioDecoderWrapper | PcmAudioDecoderWrapper;
if ((PCM_AUDIO_CODECS as readonly string[]).includes(decoderConfig.codec)) {
decoder = new PcmAudioDecoderWrapper(
sample => cursor.onDecoderSample(sample),
error => cursor.onDecoderError(error),
sample => this.onDecoderSample(sample),
error => this.onDecoderError(error),
decoderConfig,
);
} else {
decoder = new AudioDecoderWrapper(
sample => cursor.onDecoderSample(sample),
error => cursor.onDecoderError(error),
sample => this.onDecoderSample(sample),
error => this.onDecoderError(error),
codec,
decoderConfig,
);
}
decoder.onDequeue = () => cursor.onDecoderDequeue();
decoder.onDequeue = () => this.onDecoderDequeue();
const cursor: AudioSampleCursor<T> = new AudioSampleCursor(reader, decoder, options);
return cursor;
return decoder;
}
}
+60 -37
View File
@@ -20,7 +20,7 @@ test('Sample cursor seeking', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
await using cursor = await VideoSampleCursor.init(reader);
await using cursor = new VideoSampleCursor(reader);
cursor.debugInfo.enabled = true;
expect(cursor.current).toBe(null);
@@ -155,7 +155,7 @@ test('Sample cursor advancing', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor = await VideoSampleCursor.init(reader);
const cursor = new VideoSampleCursor(reader);
cursor.debugInfo.enabled = true;
expect(cursor.current).toBe(null);
@@ -250,7 +250,7 @@ test('Sample cursor advancing, cold start', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor = await VideoSampleCursor.init(reader);
const cursor = new VideoSampleCursor(reader);
cursor.debugInfo.enabled = true;
const firstSample = (await cursor.next())!;
@@ -270,7 +270,7 @@ test('Sample cursor advancing, cold start', async () => {
// Ensure the calls were serialized correctly
expect(cursor.debugInfo.seekPackets.map(x => x?.timestamp ?? null)).toEqual([0, null, 0, 2]);
const cursor2 = await VideoSampleCursor.init(reader);
const cursor2 = new VideoSampleCursor(reader);
for await (const sample of cursor2) {
expect(sample.timestamp).toBe(0);
break;
@@ -286,6 +286,28 @@ test('Sample cursor advancing, cold start', async () => {
expect(VideoSample._openSampleCount).toBe(0);
});
test('Decoder setup error', async () => {
using input = new Input({
source: new UrlSource('/trim-buck-bunny.mov'),
formats: ALL_FORMATS,
});
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
await using cursor1 = new VideoSampleCursor(reader);
cursor1.debugInfo.enabled = true;
cursor1.debugInfo.throwInDecoderInit = true;
await expect(cursor1.seekToFirst()).rejects.toThrow('Fake decoder init error');
expect(cursor1.closed).toBe(true);
const cursor2 = new VideoSampleCursor(reader);
cursor2.debugInfo.enabled = true;
cursor2.debugInfo.throwInDecoderInit = true;
await cursor2.close();
});
test('Decoder pump error handling', async () => {
using input = new Input({
source: new UrlSource('/trim-buck-bunny.mov'),
@@ -294,7 +316,7 @@ test('Decoder pump error handling', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor = await VideoSampleCursor.init(reader);
const cursor = new VideoSampleCursor(reader);
cursor.debugInfo.enabled = true;
cursor.debugInfo.throwInPump = true;
@@ -318,14 +340,14 @@ test('Decoder errors', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor1 = await VideoSampleCursor.init(reader);
const cursor1 = new VideoSampleCursor(reader);
cursor1.debugInfo.enabled = true;
cursor1.debugInfo.throwDecoderError = true;
await expect(cursor1.seekToFirst()).rejects.toThrow('Fake decoder error');
expect(cursor1.closed).toBe(true);
const cursor2 = await VideoSampleCursor.init(reader);
const cursor2 = new VideoSampleCursor(reader);
cursor2.debugInfo.enabled = true;
await cursor2.seekToFirst();
@@ -346,7 +368,7 @@ test('Use after close', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor1 = await VideoSampleCursor.init(reader);
const cursor1 = new VideoSampleCursor(reader);
expect(cursor1.closed).toBe(false);
await cursor1.close();
@@ -357,7 +379,7 @@ test('Use after close', async () => {
await expect(async () => await cursor1.seekToKey(0)).rejects.toThrow('cursor has been closed');
await expect(async () => await cursor1.next()).rejects.toThrow('cursor has been closed');
const cursor2 = await VideoSampleCursor.init(reader);
const cursor2 = new VideoSampleCursor(reader);
const commands2 = [
cursor2.seekToFirst(),
cursor2.next(),
@@ -365,12 +387,12 @@ test('Use after close', async () => {
cursor2.seekTo(1),
];
await expect(commands2[0]!).resolves.toBeInstanceOf(VideoSample);
await expect(commands2[1]!).resolves.toBeInstanceOf(VideoSample);
await expect(commands2[2]!).resolves.toBeUndefined();
await expect(commands2[3]!).rejects.toThrow('cursor has been closed');
await expect(commands2[0]).resolves.toBeInstanceOf(VideoSample);
await expect(commands2[1]).resolves.toBeInstanceOf(VideoSample);
await expect(commands2[2]).resolves.toBeUndefined();
await expect(commands2[3]).rejects.toThrow('cursor has been closed');
const cursor3 = await VideoSampleCursor.init(reader);
const cursor3 = new VideoSampleCursor(reader);
const commands3 = [
cursor3.seekToFirst(),
cursor3.next(),
@@ -378,10 +400,10 @@ test('Use after close', async () => {
cursor3.next(),
];
await expect(commands3[0]!).resolves.toBeInstanceOf(VideoSample);
await expect(commands3[1]!).resolves.toBeInstanceOf(VideoSample);
await expect(commands3[2]!).resolves.toBeUndefined();
await expect(commands3[3]!).rejects.toThrow('cursor has been closed');
await expect(commands3[0]).resolves.toBeInstanceOf(VideoSample);
await expect(commands3[1]).resolves.toBeInstanceOf(VideoSample);
await expect(commands3[2]).resolves.toBeUndefined();
await expect(commands3[3]).rejects.toThrow('cursor has been closed');
expect(VideoSample._openSampleCount).toBe(0);
});
@@ -396,7 +418,7 @@ test('Command queuing', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor0 = await VideoSampleCursor.init(reader);
const cursor0 = new VideoSampleCursor(reader);
cursor0.debugInfo.enabled = true;
const commands0 = [
@@ -406,9 +428,12 @@ test('Command queuing', async () => {
];
expect(commands0.every(x => x instanceof Promise)).toBe(true);
// eslint-disable-next-line @typescript-eslint/await-thenable
await Promise.all(commands0);
expect(cursor0.debugInfo.pumpsStarted).toBe(1);
const cursor1 = await VideoSampleCursor.init(reader);
const cursor1 = new VideoSampleCursor(reader);
cursor1.debugInfo.enabled = true;
const commands1 = [
@@ -422,7 +447,7 @@ test('Command queuing', async () => {
expect(results1[0]!.timestamp).toBe(0);
expect(cursor1.debugInfo.decodedPackets.map(x => x.timestamp)).toEqual([0]);
const cursor2 = await VideoSampleCursor.init(reader);
const cursor2 = new VideoSampleCursor(reader);
cursor2.debugInfo.enabled = true;
const commands2 = [
@@ -448,7 +473,7 @@ test('Command queuing', async () => {
0, 1, 2, 3, 4, 5,
]);
const cursor3 = await VideoSampleCursor.init(reader);
const cursor3 = new VideoSampleCursor(reader);
cursor3.debugInfo.enabled = true;
const commands3 = [
@@ -474,7 +499,7 @@ test('Command queuing', async () => {
expect(cursor3.debugInfo.decodedPackets.every(x => x.timestamp <= 0.5)).toBe(true);
expect(cursor3.debugInfo.pumpsStarted).toBe(1);
const cursor4 = await VideoSampleCursor.init(reader);
const cursor4 = new VideoSampleCursor(reader);
cursor4.debugInfo.enabled = true;
const commands4 = [
@@ -494,7 +519,7 @@ test('Command queuing', async () => {
expect(results4[2]!.timestamp).toBeGreaterThan(results4[1]!.timestamp);
expect(cursor4.debugInfo.decodedPackets.length).toBeGreaterThan(3); // Because .next() goes into "sequential mode"
const cursor5 = await VideoSampleCursor.init(reader);
const cursor5 = new VideoSampleCursor(reader);
cursor5.debugInfo.enabled = true;
const commands5 = [
@@ -531,7 +556,7 @@ test('Command queuing', async () => {
expect(cursor5.debugInfo.pumpsStarted).toBe(1);
const cursor6 = await VideoSampleCursor.init(reader);
const cursor6 = new VideoSampleCursor(reader);
cursor6.debugInfo.enabled = true;
const commands6 = [
@@ -556,7 +581,7 @@ test('Command queuing', async () => {
expect(cursor6.debugInfo.pumpsStarted).toBe(3);
const cursor7 = await VideoSampleCursor.init(reader);
const cursor7 = new VideoSampleCursor(reader);
const commands7 = [
cursor7.close(),
cursor7.close(),
@@ -565,7 +590,7 @@ test('Command queuing', async () => {
await Promise.all(commands7);
const cursor8 = await VideoSampleCursor.init(reader, { autoClose: false });
const cursor8 = new VideoSampleCursor(reader, { autoClose: false });
const firstSample = await cursor8.seekToFirst();
firstSample!.close();
@@ -596,7 +621,7 @@ test('Automatic cursor disposal', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor = await VideoSampleCursor.init(reader);
const cursor = new VideoSampleCursor(reader);
await cursor.seekToFirst();
// No cursor.close() here, but the disposed Input closes the cursor
@@ -616,7 +641,7 @@ test('Video with stubborn first sample emit', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
await using cursor = await VideoSampleCursor.init(reader);
await using cursor = new VideoSampleCursor(reader);
const firstSample = (await cursor.seekToFirst())!;
expect(firstSample).not.toBe(null);
@@ -631,7 +656,7 @@ test('AudioSampleCursor', async () => {
const audioTrack = (await input.getPrimaryAudioTrack())!;
const reader = new PacketReader(audioTrack);
const cursor = await AudioSampleCursor.init(reader);
const cursor = new AudioSampleCursor(reader);
cursor.debugInfo.enabled = true;
const firstSample = (await cursor.seekToFirst())!;
@@ -679,6 +704,7 @@ test('AudioSampleCursor', async () => {
cursor.seekTo(0.15),
cursor.seekTo(0.2),
];
// eslint-disable-next-line @typescript-eslint/await-thenable
const result = await Promise.all(commands);
@@ -706,7 +732,7 @@ test('Sample mapping', async () => {
const reader = new PacketReader(videoTrack);
let callCount = 0;
await using cursor = await VideoSampleCursor.init(reader, {
await using cursor = new VideoSampleCursor(reader, {
transform: (sample) => {
callCount++;
@@ -743,7 +769,7 @@ test('Canvas transformer', async () => {
const videoTrack = (await input.getPrimaryVideoTrack())!;
const reader = new PacketReader(videoTrack);
const cursor1 = await VideoSampleCursor.init(reader, {
const cursor1 = new VideoSampleCursor(reader, {
transform: canvasTransformer(),
});
@@ -761,7 +787,7 @@ test('Canvas transformer', async () => {
await cursor1.close();
const cursor2 = await VideoSampleCursor.init(reader, {
const cursor2 = new VideoSampleCursor(reader, {
transform: canvasTransformer({
width: 320,
poolSize: 2,
@@ -797,7 +823,7 @@ test('AudioBuffer transformer', async () => {
const audioTrack = (await input.getPrimaryAudioTrack())!;
const reader = new PacketReader(audioTrack);
const cursor = await AudioSampleCursor.init(reader, {
const cursor = new AudioSampleCursor(reader, {
transform: audioBufferTransformer(),
});
@@ -817,6 +843,3 @@ test('AudioBuffer transformer', async () => {
expect(AudioSample._openSampleCount).toBe(0);
});
// TODO:
// - Then, clean up the cursor.ts code