diff --git a/src/cursors.ts b/src/cursors.ts index bf57921..9f6ee81 100644 --- a/src/cursors.ts +++ b/src/cursors.ts @@ -1,11 +1,14 @@ // Two parallelism modes: Cancel and queue (right?) // Useful in packet context? Or only for sample? -import { InputTrack, InputVideoTrack } from './input-track'; -import { PacketRetrievalOptions, validatePacketRetrievalOptions, validateTimestamp, VideoDecoderWrapper } from './media-sink'; -import { assert, assertNever, AsyncMutex, AsyncMutex2, AsyncMutex3, AsyncMutexLock, CallSerializer2, defer, insertSorted, last, MaybePromise, promiseWithResolvers, ResultValue, Yo } from './misc'; +import { PCM_AUDIO_CODECS } from './codec'; +import { InputAudioTrack, InputTrack, InputVideoTrack } from './input-track'; +import { AudioDecoderWrapper, DecoderWrapper, PacketRetrievalOptions, PcmAudioDecoderWrapper, validatePacketRetrievalOptions, validateTimestamp, VideoDecoderWrapper } from './media-sink'; +import { assert, AsyncMutex3, AsyncMutexLock, CallSerializer2, defer, insertSorted, isFirefox, last, MaybePromise, polyfillSymbolDispose, promiseWithResolvers, ResultValue, Rotation, Yo } from './misc'; import { EncodedPacket } from './packet'; -import { VideoSample } from './sample'; +import { AudioSample, clampCropRectangle, CropRectangle, validateCropRectangle, VideoSample } from './sample'; + +polyfillSymbolDispose(); export class PacketReader { track: T; @@ -332,335 +335,38 @@ export class PacketCursor { } } -/* -export class PacketCursor { - track: InputTrack; - current: EncodedPacket | null = null; - nextIsFirst = true; - callSerializer = new CallSerializer2(); - - constructor(track: InputTrack) { - if (!(track instanceof InputTrack)) { - throw new TypeError('track must be an InputTrack.'); - } - - this.track = track; - } - - peekAtStart(options: PacketRetrievalOptions = {}): MaybePromise { - validatePacketRetrievalOptions(options); - - const result = new ResultValue(); - const promise = this.track._backing.getFirstPacket(result, options); - - if (result.pending) { - return (promise as Promise).then(() => result.value); - } else { - return result.value; - } - } - - seekToStartDirect(options: PacketRetrievalOptions): MaybePromise { - const result = new ResultValue(); - const promise = this.track._backing.getFirstPacket(result, options); - - const onPacket = () => { - this.nextIsFirst = false; - return this.current = result.value; - }; - - if (result.pending) { - return (promise as Promise).then(onPacket); - } else { - return onPacket(); - } - } - - seekToStart(options: PacketRetrievalOptions = {}): MaybePromise { - validatePacketRetrievalOptions(options); - - return this.callSerializer.call(() => this.seekToStartDirect(options)); - } - - peekAt(timestamp: number, options: PacketRetrievalOptions = {}): MaybePromise { - validateTimestamp(timestamp); - validatePacketRetrievalOptions(options); - - const result = new ResultValue(); - const promise = this.track._backing.getPacket(result, timestamp, options); - - if (result.pending) { - return (promise as Promise).then(() => result.value); - } else { - return result.value; - } - } - - seekTo(timestamp: number, options: PacketRetrievalOptions = {}): MaybePromise { - validateTimestamp(timestamp); - validatePacketRetrievalOptions(options); - - return this.callSerializer.call(() => { - const result = new ResultValue(); - const promise = this.track._backing.getPacket(result, timestamp, options); - - const onPacket = () => { - this.nextIsFirst = !result.value; - return this.current = result.value; - }; - - if (result.pending) { - return (promise as Promise).then(onPacket); - } else { - return onPacket(); - } - }); - } - - peekKeyAt(timestamp: number, options: PacketRetrievalOptions = {}): MaybePromise { - validateTimestamp(timestamp); - validatePacketRetrievalOptions(options); - - const result = new ResultValue(); - const promise = this.track._backing.getKeyPacket(result, timestamp, options); - - if (result.pending) { - return (promise as Promise).then(() => result.value); - } else { - return result.value; - } - } - - seekToKey(timestamp: number, options: PacketRetrievalOptions = {}): MaybePromise { - validateTimestamp(timestamp); - validatePacketRetrievalOptions(options); - - return this.callSerializer.call(() => { - const result = new ResultValue(); - const promise = this.track._backing.getKeyPacket(result, timestamp, options); - - const onPacket = () => { - this.nextIsFirst = !result.value; - return this.current = result.value; - }; - - if (result.pending) { - return (promise as Promise).then(onPacket); - } else { - return onPacket(); - } - }); - } - - peekNext(from?: EncodedPacket | null, options: PacketRetrievalOptions = {}): MaybePromise { - if (from != null && !(from instanceof EncodedPacket)) { - throw new TypeError('from, when provided, must be an EncodedPacket or null.'); - } - validatePacketRetrievalOptions(options); - - const run = () => { - if (from === null) { - if (this.nextIsFirst) { - return this.peekAtStart(); - } else { - return null; - } - } - - assert(from); - - const result = new ResultValue(); - const promise = this.track._backing.getNextPacket(result, from, options); - - if (result.pending) { - return (promise as Promise).then(() => result.value); - } else { - return result.value; - } - }; - - if (from === undefined) { - return this.callSerializer.call(() => { - from = this.current; - return run(); - }); - } else { - return run(); - } - } - - next(options: PacketRetrievalOptions = {}): MaybePromise { - validatePacketRetrievalOptions(options); - - return this.callSerializer.call(() => { - if (this.nextIsFirst) { - return this.seekToStartDirect(options); - } - - if (!this.current) { - return null; - } - - const result = new ResultValue(); - const promise = this.track._backing.getNextPacket(result, this.current, options); - - const onPacket = () => { - return this.current = result.value; - }; - - if (result.pending) { - return (promise as Promise).then(onPacket); - } else { - return onPacket(); - } - }); - } - - peekNextKey(from?: EncodedPacket | null, options: PacketRetrievalOptions = {}): MaybePromise { - if (from != null && !(from instanceof EncodedPacket)) { - throw new TypeError('from, when provided, must be an EncodedPacket or null.'); - } - validatePacketRetrievalOptions(options); - - const run = () => { - if (from === null) { - if (this.nextIsFirst) { - return this.peekAtStart(); - } else { - return null; - } - } - - assert(from); - - const result = new ResultValue(); - const promise = this.track._backing.getNextKeyPacket(result, from, options); - - if (result.pending) { - return (promise as Promise).then(() => result.value); - } else { - return result.value; - } - }; - - if (from === undefined) { - return this.callSerializer.call(() => { - from = this.current; - return run(); - }); - } else { - return run(); - } - } - - nextKey(options: PacketRetrievalOptions = {}): MaybePromise { - return this.callSerializer.call(() => { - if (this.nextIsFirst) { - return this.seekToStartDirect(options); - } - - if (!this.current) { - return null; - } - - const result = new ResultValue(); - const promise = this.track._backing.getNextKeyPacket(result, this.current, options); - - if (result.pending) { - return (promise as Promise).then(() => { - return this.current = result.value; - }); - } else { - return this.current = result.value; - } - }); - } - - async iterate( - callback: (packet: EncodedPacket, stop: () => void) => MaybePromise, - options: PacketRetrievalOptions = {}, - ) { - let stopped = false; - const stop = () => stopped = true; - - const donePromise = this.callSerializer.done(); - if (donePromise) await donePromise; - - while (true) { - if (this.current) { - const result = callback(this.current, stop); - if (result instanceof Promise) await result; - } - - if (stopped) { - break; - } - - const result = this.next(); - if (result instanceof Promise) await result; - - if (!this.current) { - break; - } - } - } - - // eslint-disable-next-line @stylistic/generator-star-spacing - async *[Symbol.asyncIterator]() { - const donePromise = this.callSerializer.done(); - if (donePromise) await donePromise; - - while (true) { - if (this.current) { - yield this.current; - } - - const result = this.next(); - if (result instanceof Promise) await result; - - if (!this.current) { - break; - } - } - } - - waitUntilIdle() { - return this.callSerializer.done(); - } -} -*/ - -/* -for (const timestamp of timestamps) { - const thumbnail = await cursor.seekTo(timestamp); - console.log(thumbnail); -} -*/ - -type PendingRequest = { +type PendingRequest = { timestamp: number; - promise: Promise; - resolve: (sample: VideoSample | null) => void; + promise: Promise; + resolve: (sample: T | null) => void; reject: (error: unknown) => void; - successor: PendingRequest | null; + successor: PendingRequest | null; }; -type VideoSampleCursorOptions = { +type SampleTransformer = (sample: Sample) => MaybePromise; + +type SampleCursorOptions = { autoClose?: boolean; + transform?: SampleTransformer; }; -export class VideoSampleCursor { - packetReader: PacketReader; +export abstract class SampleCursor< + Sample extends VideoSample | AudioSample, + TransformedSample = Sample, +> implements AsyncDisposable { + packetReader: PacketReader; packetCursor: PacketCursor; - options: VideoSampleCursorOptions; + options: SampleCursorOptions; + transform: SampleTransformer; autoClose: boolean; pumpRunning = false; - decoder: VideoDecoderWrapper; - current: VideoSample | null = null; - sampleQueue: VideoSample[] = []; + decoder: DecoderWrapper; + currentRaw: Sample | null = null; + current: TransformedSample | null = null; + sampleQueue: Sample[] = []; queueDequeue = promiseWithResolvers(); - pendingRequests: PendingRequest[] = []; - lastPendingRequest: PendingRequest | null = null; + pendingRequests: PendingRequest[] = []; + lastPendingRequest: PendingRequest | null = null; predictedRequests = 0; nextIsFirst = true; @@ -672,7 +378,11 @@ export class VideoSampleCursor { maxDecodedSequenceNumber = -1; pumpTarget: EncodedPacket | null = null; pumpMutex = new AsyncMutex3(); - closed = false; + _closed = false; + otherMutex = new AsyncMutex3(); // TODO: THIS IS STILL A BUGGED MUTEX! ASYNC MUTEX 4 + + error: unknown = null; + errorSet = false; debugInfo = { enabled: false, @@ -680,102 +390,108 @@ export class VideoSampleCursor { seekPackets: [] as (EncodedPacket | null)[], decodedPackets: [] as EncodedPacket[], throwInPump: false, + throwDecoderError: false, }; - private constructor( - reader: PacketReader, - decoder: VideoDecoderWrapper, - options: VideoSampleCursorOptions, + get closed(): boolean { + return this._closed; + } + + protected constructor( + reader: PacketReader, + decoder: DecoderWrapper, + options: SampleCursorOptions, ) { 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); } - static async init(reader: PacketReader, options: VideoSampleCursorOptions = {}) { - if (!(reader instanceof PacketReader) || !(reader.track instanceof InputVideoTrack)) { - throw new TypeError('reader must be a PacketReader for an InputVideoTrack.'); - } + async onDecoderSample(sample: Sample) { + try { + if (this.debugInfo.enabled && this.debugInfo.throwDecoderError) { + sample.close(); + return this.onDecoderError(new Error('Fake decoder error!')); + } - const track = reader.track; + const mutexPromise = this.otherMutex.request(); + if (mutexPromise) await mutexPromise; + using _ = this.otherMutex.lock(); - if (!(await track.canDecode())) { - throw new Error( - 'This video track cannot be decoded by this browser. Make sure to check decodability before using' - + ' a track.', - ); - } + while (this.decodedTimestamps.length > 0 && this.decodedTimestamps[0]! <= sample.timestamp) { + this.decodedTimestamps.shift(); + } - const decoderConfig = await track.getDecoderConfig(); - assert(decoderConfig); - assert(track.codec); - - const decoder = new VideoDecoderWrapper( - (sample) => { - while (cursor.decodedTimestamps.length > 0 && cursor.decodedTimestamps[0]! <= sample.timestamp) { - cursor.decodedTimestamps.shift(); - } - - if (cursor.pendingRequests.length === 0) { - if (cursor.pumpStopQueued) { - sample.close(); - } else { - cursor.sampleQueue.push(sample); - } + if (this.pendingRequests.length === 0) { + if (this.pumpStopQueued) { + sample.close(); } else { - let given = false; + this.sampleQueue.push(sample); + } + } else { + let given = false; + let transformed: TransformedSample | null = null; - for (let i = 0; i < cursor.pendingRequests.length; i++) { - const request = cursor.pendingRequests[i]!; - if (request.timestamp > sample.timestamp) { - break; - } - - request.resolve(sample); - cursor.pendingRequests.splice(i--, 1); - cursor.setCurrent(sample); - given = true; - - if (request.successor) { - cursor.pendingRequests.unshift(request.successor); - i++; - } + for (let i = 0; i < this.pendingRequests.length; i++) { + const request = this.pendingRequests[i]!; + if (request.timestamp > sample.timestamp) { + break; } - if (!given) { - sample.close(); + if (!transformed) { + let result = this.transform(sample); + if (result instanceof Promise) result = await result; + + transformed = result; + } + + request.resolve(transformed); + this.pendingRequests.splice(i--, 1); + this.setCurrent(sample, transformed); + given = true; + + if (request.successor) { + this.pendingRequests.unshift(request.successor); + i++; } } - cursor.queueDequeue.resolve(); - cursor.queueDequeue = promiseWithResolvers(); - }, - (error) => { - // TODO THIS TODO THIS TODO THIS TODO THIS - console.error(error); - }, - track.codec, - decoderConfig, - track.rotation, - track.timeResolution, - ); + if (!given) { + sample.close(); + } + } - decoder.onDequeue = () => { - cursor.queueDequeue.resolve(); - cursor.queueDequeue = promiseWithResolvers(); - }; - - const cursor = new VideoSampleCursor(reader, decoder, options); - return cursor; + this.queueDequeue.resolve(); + this.queueDequeue = promiseWithResolvers(); + } catch (error) { + await this.closeWithError(error); + } } - setCurrent(newCurrent: VideoSample | null) { - if (this.autoClose && this.current && this.current !== newCurrent) { - this.current.close(); + async onDecoderError(error: unknown) { + const mutexPromise = this.otherMutex.request(); + if (mutexPromise) await mutexPromise; + using _ = this.otherMutex.lock(); + + await this.closeWithError(error); + } + + onDecoderDequeue() { + this.queueDequeue.resolve(); + this.queueDequeue = promiseWithResolvers(); + } + + setCurrent(newCurrentRaw: Sample | null, newCurrent: TransformedSample | null) { + if (this.autoClose && this.currentRaw && this.currentRaw !== newCurrentRaw) { + this.currentRaw.close(); } + this.currentRaw = newCurrentRaw; this.current = newCurrent; } @@ -787,12 +503,16 @@ export class VideoSampleCursor { _ensureNotClosed() { if (this.closed) { - throw new Error('This cursor has been closed and can no longer be used.'); + if (this.errorSet) { + throw this.error; + } else { + throw new Error('This cursor has been closed and can no longer be used.'); + } } } async _seekToPacket( - res: ResultValue, + res: ResultValue, targetPacketPromise: MaybePromise, lock?: AsyncMutexLock, ): Promise { @@ -800,12 +520,12 @@ export class VideoSampleCursor { if (!lock) { const mutexPromise = this.pumpMutex.request(); - if (mutexPromise) { - await mutexPromise; - } + if (mutexPromise) await mutexPromise; lock = this.pumpMutex.lock(); } + this._ensureNotClosed(); + using deferred = defer(() => { lock?.release(); this.predictedRequests--; @@ -822,11 +542,11 @@ export class VideoSampleCursor { this.nextIsFirst = !targetPacket; if (!targetPacket) { - this.setCurrent(null); + this.setCurrent(null, null); return res.set(null); } - if (this.current?.timestamp === targetPacket.timestamp) { + if (this.currentRaw?.timestamp === targetPacket.timestamp) { return res.set(this.current); } @@ -836,16 +556,19 @@ export class VideoSampleCursor { } else { while (this.sampleQueue.length > 0) { - const nextSample = this.sampleQueue[0]!; - if (targetPacket.timestamp <= nextSample.timestamp) { - this.setCurrent(nextSample); - return res.set(nextSample); - } - - this.sampleQueue.shift(); + const nextSample = this.sampleQueue.shift()!; this.queueDequeue.resolve(); this.queueDequeue = promiseWithResolvers(); - nextSample.close(); + + if (targetPacket.timestamp <= nextSample.timestamp) { + let transformed = this.transform(nextSample); + if (transformed instanceof Promise) transformed = await transformed; + + this.setCurrent(nextSample, transformed); + return res.set(transformed); + } else { + nextSample.close(); + } } if (this.pumpTarget) { @@ -892,8 +615,10 @@ export class VideoSampleCursor { void this.runPump(); } - const request = promiseWithResolvers(); - const pendingRequest: PendingRequest = { + this._ensureNotClosed(); + + const request = promiseWithResolvers(); + const pendingRequest: PendingRequest = { timestamp: targetPacket.timestamp, promise: request.promise, resolve: request.resolve, @@ -912,50 +637,73 @@ export class VideoSampleCursor { return res.set(await request.promise); } - seekToFirst(): MaybePromise { + seekToFirst(): MaybePromise { this._ensureNotClosed(); - const result = new ResultValue(); - const promise = this._seekToPacket(result, this.packetReader.readFirst()); + try { + const result = new ResultValue(); + const promise = this._seekToPacket(result, this.packetReader.readFirst()); - if (result.pending) { - return promise.then(() => result.value); - } else { - return result.value; + if (result.pending) { + return promise + .then(() => result.value) + .catch(this.closeWithErrorAndThrow.bind(this)); + } else { + return result.value; + } + } catch (error) { + this.closeWithErrorAndThrow(error); } } - seekTo(timestamp: number): MaybePromise { + seekTo(timestamp: number): MaybePromise { this._ensureNotClosed(); - const result = new ResultValue(); - const promise = this._seekToPacket(result, this.packetReader.readAt(timestamp)); + try { + const result = new ResultValue(); + const promise = this._seekToPacket(result, this.packetReader.readAt(timestamp)); - if (result.pending) { - return promise.then(() => result.value); - } else { - return result.value; + if (result.pending) { + return promise + .then(() => result.value) + .catch(this.closeWithErrorAndThrow.bind(this)); + } else { + return result.value; + } + } catch (error) { + this.closeWithErrorAndThrow(error); } } - seekToKey(timestamp: number): MaybePromise { + seekToKey(timestamp: number): MaybePromise { this._ensureNotClosed(); - const result = new ResultValue(); - const promise = this._seekToPacket(result, this.packetReader.readKeyAt(timestamp, { verifyKeyPackets: true })); + try { + const result = new ResultValue(); + const promise = this._seekToPacket( + result, + this.packetReader.readKeyAt(timestamp, { verifyKeyPackets: true }), + ); - if (result.pending) { - return promise.then(() => result.value); - } else { - return result.value; + if (result.pending) { + return promise + .then(() => result.value) + .catch(this.closeWithErrorAndThrow.bind(this)); + } else { + return result.value; + } + } catch (error) { + this.closeWithErrorAndThrow(error); } } - async _nextInternal(res: ResultValue): Promise { + async _nextInternal(res: ResultValue): Promise { const mutexPromise = this.pumpMutex.request(); if (mutexPromise) await mutexPromise; using lock = this.pumpMutex.lock(); + this._ensureNotClosed(); + if (this.nextIsFirst) { return await this._seekToPacket(res, this.packetReader.readFirst(), lock); } @@ -965,19 +713,24 @@ export class VideoSampleCursor { this.queueDequeue.resolve(); this.queueDequeue = promiseWithResolvers(); - this.setCurrent(nextSample); - return res.set(nextSample); + let transformed = this.transform(nextSample); + if (transformed instanceof Promise) transformed = await transformed; + + this.setCurrent(nextSample, transformed); + return res.set(transformed); } if (!this.pumpRunning) { - this.setCurrent(null); + this.setCurrent(null, null); return res.set(null); // None more after this, boy } assert(this.lastPendingRequest); - const request = promiseWithResolvers(); - const pendingRequest: PendingRequest = { + this._ensureNotClosed(); + + const request = promiseWithResolvers(); + const pendingRequest: PendingRequest = { timestamp: -Infinity, promise: request.promise, resolve: request.resolve, @@ -1005,22 +758,30 @@ export class VideoSampleCursor { return res.set(await request.promise); } - next(): MaybePromise { + next(): MaybePromise { this._ensureNotClosed(); - const result = new ResultValue(); - const promise = this._nextInternal(result); + try { + const result = new ResultValue(); + const promise = this._nextInternal(result); - if (result.pending) { - return promise.then(() => result.value); - } else { - return result.value; + if (result.pending) { + return promise + .then(() => result.value) + .catch(this.closeWithErrorAndThrow.bind(this)); + } else { + return result.value; + } + } catch (error) { + this.closeWithErrorAndThrow(error); } } async iterate( - callback: (packet: VideoSample, stop: () => void) => MaybePromise, + callback: (sample: TransformedSample, stop: () => void) => MaybePromise, ) { + this._ensureNotClosed(); + let stopped = false; const stop = () => stopped = true; @@ -1048,6 +809,8 @@ export class VideoSampleCursor { // eslint-disable-next-line @stylistic/generator-star-spacing async *[Symbol.asyncIterator]() { + this._ensureNotClosed(); + const waitPromise = this.waitUntilIdle(); if (waitPromise) await waitPromise; @@ -1085,21 +848,27 @@ export class VideoSampleCursor { close() { return this.closePromise ??= (async () => { this.predictedRequests++; + this.packetReader.track.input._openSampleCursors.delete(this); const mutexPromise = this.pumpMutex.request(); if (mutexPromise) await mutexPromise; + using _ = this.pumpMutex.lock(); - this.closed = true; + this._closed = true; if (this.pumpRunning) { await this.stopPump(); } - this.setCurrent(null); + this.setCurrent(null, null); this.decoder.close(); })(); } + [Symbol.asyncDispose]() { + return this.close(); + } + async stopPump() { assert(this.pumpRunning); @@ -1129,7 +898,7 @@ export class VideoSampleCursor { ) ) { if (this.debugInfo.enabled && this.debugInfo.throwInPump) { - throw new Error('Throwing artificially.'); + throw new Error('Fake pump error!'); } if ( @@ -1158,11 +927,11 @@ export class VideoSampleCursor { await this.packetCursor.next(); } - if (this.pendingRequests.length > 0 || !this.closed) { + if (this.pendingRequests.length > 0 || !this._closed) { await this.decoder.flush(); } - const uh = (request: PendingRequest) => { + const uh = (request: PendingRequest) => { request.resolve(null); if (request.successor) { uh(request.successor); @@ -1170,27 +939,14 @@ export class VideoSampleCursor { }; this.pendingRequests.forEach(uh); - this.setCurrent(null); + this.setCurrent(null, null); } catch (error) { - if (!this.decoder.closed) { - if (this.pendingRequests.length > 0 || !this.closed) { - await this.decoder.flush(); - } + if (!this.decoder.closed && this.pendingRequests.length > 0) { + await this.decoder.flush(); } - const uh = (request: PendingRequest) => { - request.reject(error); - if (request.successor) { - uh(request.successor); - } - }; - - if (this.pendingRequests.length > 0) { - this.pendingRequests.forEach(uh); - this.setCurrent(null); - } else { - throw error; // To make sure it isn't lost - } + this.pumpRunning = false; // So that close() doesn't attempt to stop the pump + void this.closeWithError(error); } finally { for (const sample of this.sampleQueue) { sample.close(); @@ -1208,4 +964,338 @@ export class VideoSampleCursor { this.pumpStopQueued = false; } } + + closeWithError(error: unknown) { + if (this._closed) { + return; + } + + this._closed = true; + this.error = error; + this.errorSet = true; + + const uh = (request: PendingRequest) => { + request.reject(error); + if (request.successor) { + uh(request.successor); + } + }; + this.pendingRequests.forEach(uh); + + return this.close(); + } + + closeWithErrorAndThrow(error: unknown): never { + void this.closeWithError(error); + throw error; + } } + +export class VideoSampleCursor extends SampleCursor { + static async init( + reader: PacketReader, + options: SampleCursorOptions = {}, + ) { + if (!(reader instanceof PacketReader) || !(reader.track instanceof InputVideoTrack)) { + throw new TypeError('reader must be a PacketReader for an InputVideoTrack.'); + } + + const track = reader.track; + + if (!(await track.canDecode())) { + throw new Error( + 'This video track cannot be decoded by this browser. Make sure to check decodability before using' + + ' a track.', + ); + } + + const decoderConfig = await track.getDecoderConfig(); + assert(decoderConfig); + assert(track.codec); + + const decoder = new VideoDecoderWrapper( + sample => cursor.onDecoderSample(sample), + error => cursor.onDecoderError(error), + track.codec, + decoderConfig, + track.rotation, + track.timeResolution, + ); + + decoder.onDequeue = () => cursor.onDecoderDequeue(); + + const cursor: VideoSampleCursor = new VideoSampleCursor(reader, decoder, options); + return cursor; + } +} + +export class AudioSampleCursor extends SampleCursor { + static async init( + reader: PacketReader, + options: SampleCursorOptions = {}, + ) { + if (!(reader instanceof PacketReader) || !(reader.track instanceof InputAudioTrack)) { + throw new TypeError('reader must be a PacketReader for an InputAudioTrack.'); + } + + const track = reader.track; + + if (!(await track.canDecode())) { + throw new Error( + 'This audio track cannot be decoded by this browser. Make sure to check decodability before using' + + ' a track.', + ); + } + + const codec = track.codec; + const decoderConfig = await track.getDecoderConfig(); + assert(codec && decoderConfig); + + 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), + decoderConfig, + ); + } else { + decoder = new AudioDecoderWrapper( + sample => cursor.onDecoderSample(sample), + error => cursor.onDecoderError(error), + codec, + decoderConfig, + ); + } + + decoder.onDequeue = () => cursor.onDecoderDequeue(); + + const cursor: AudioSampleCursor = new AudioSampleCursor(reader, decoder, options); + return cursor; + } +} + +/** + * A canvas with additional timing information (timestamp & duration). + * @public + */ +export type WrappedCanvas = { + /** A canvas element or offscreen canvas. */ + canvas: HTMLCanvasElement | OffscreenCanvas; + /** The timestamp of the corresponding video sample, in seconds. */ + timestamp: number; + /** The duration of the corresponding video sample, in seconds. */ + duration: number; +}; + +/** + * Options for constructing a canvas transformer. + * @public + */ +export type CanvasTransformerOptions = { + /** + * Whether the output canvases should have transparency instead of a black background. Defaults to `false`. Set + * this to `true` when using this sink to read transparent videos. + */ + alpha?: boolean; + /** + * The width of the output canvas in pixels, defaulting to the display width of the video track. If height is not + * set, it will be deduced automatically based on aspect ratio. + */ + width?: number; + /** + * The height of the output canvas in pixels, defaulting to the display height of the video track. If width is not + * set, it will be deduced automatically based on aspect ratio. + */ + height?: number; + /** + * The fitting algorithm in case both width and height are set. + * + * - `'fill'` will stretch the image to fill the entire box, potentially altering aspect ratio. + * - `'contain'` will contain the entire image within the box while preserving aspect ratio. This may lead to + * letterboxing. + * - `'cover'` will scale the image until the entire box is filled, while preserving aspect ratio. + */ + fit?: 'fill' | 'contain' | 'cover'; + /** + * The clockwise rotation by which to rotate the raw video frame. Defaults to the rotation set in the file metadata. + * Rotation is applied before resizing. + */ + rotation?: Rotation; + /** + * Specifies the rectangular region of the input video to crop to. The crop region will automatically be clamped to + * the dimensions of the input video track. Cropping is performed after rotation but before resizing. + */ + crop?: CropRectangle; + /** + * When set, specifies the number of canvases in the pool. These canvases will be reused in a ring buffer / + * round-robin type fashion. This keeps the amount of allocated VRAM constant and relieves the browser from + * constantly allocating/deallocating canvases. A pool size of 0 or `undefined` disables the pool and means a new + * canvas is created each time. + */ + poolSize?: number; +}; + +export const canvasTransformer = ( + options: CanvasTransformerOptions = {}, +): SampleTransformer => { + if (options && typeof options !== 'object') { + throw new TypeError('options must be an object.'); + } + if (options.alpha !== undefined && typeof options.alpha !== 'boolean') { + throw new TypeError('options.alpha, when provided, must be a boolean.'); + } + if (options.width !== undefined && (!Number.isInteger(options.width) || options.width <= 0)) { + throw new TypeError('options.width, when defined, must be a positive integer.'); + } + if (options.height !== undefined && (!Number.isInteger(options.height) || options.height <= 0)) { + throw new TypeError('options.height, when defined, must be a positive integer.'); + } + if (options.fit !== undefined && !['fill', 'contain', 'cover'].includes(options.fit)) { + throw new TypeError('options.fit, when provided, must be one of "fill", "contain", or "cover".'); + } + if ( + options.width !== undefined + && options.height !== undefined + && options.fit === undefined + ) { + throw new TypeError( + 'When both options.width and options.height are provided, options.fit must also be provided.', + ); + } + if (options.rotation !== undefined && ![0, 90, 180, 270].includes(options.rotation)) { + throw new TypeError('options.rotation, when provided, must be 0, 90, 180 or 270.'); + } + if (options.crop !== undefined) { + validateCropRectangle(options.crop, 'options.'); + } + if ( + options.poolSize !== undefined + && (typeof options.poolSize !== 'number' || !Number.isInteger(options.poolSize) || options.poolSize < 0) + ) { + throw new TypeError('poolSize must be a non-negative integer.'); + } + + let needsSetup = true; + let alpha: boolean; + let width: number; + let height: number; + let fit: 'fill' | 'contain' | 'cover'; + let rotation: Rotation; + let crop: { left: number; top: number; width: number; height: number } | undefined; + let canvasPool: (HTMLCanvasElement | OffscreenCanvas | null)[]; + let nextCanvasIndex = 0; + + return (sample) => { + if (needsSetup) { + rotation = options.rotation ?? sample.rotation; + + const [rotatedWidth, rotatedHeight] = rotation % 180 === 0 + ? [sample.codedWidth, sample.codedHeight] + : [sample.codedHeight, sample.codedWidth]; + + crop = options.crop; + if (crop) { + clampCropRectangle(crop, rotatedWidth, rotatedHeight); + } + + [width, height] = crop + ? [crop.width, crop.height] + : [rotatedWidth, rotatedHeight]; + const originalAspectRatio = width / height; + + // If width and height aren't defined together, deduce the missing value using the aspect ratio + if (options.width !== undefined && options.height === undefined) { + width = options.width; + height = Math.round(width / originalAspectRatio); + } else if (options.width === undefined && options.height !== undefined) { + height = options.height; + width = Math.round(height * originalAspectRatio); + } else if (options.width !== undefined && options.height !== undefined) { + width = options.width; + height = options.height; + } + + alpha = options.alpha ?? false; + fit = options.fit ?? 'fill'; + canvasPool = Array.from({ length: options.poolSize ?? 0 }, () => null); + needsSetup = false; + } + + let canvas = canvasPool[nextCanvasIndex]; + let canvasIsNew = false; + + if (!canvas) { + if (typeof document !== 'undefined') { + // Prefer an HTMLCanvasElement + canvas = document.createElement('canvas'); + canvas.width = width; + canvas.height = height; + } else { + canvas = new OffscreenCanvas(width, height); + } + + if (canvasPool.length > 0) { + canvasPool[nextCanvasIndex] = canvas; + } + + canvasIsNew = true; + } + + if (canvasPool.length > 0) { + nextCanvasIndex = (nextCanvasIndex + 1) % canvasPool.length; + } + + const context = canvas.getContext('2d', { + alpha: alpha || isFirefox(), // Firefox has VideoFrame glitches with opaque canvases + }) as CanvasRenderingContext2D | OffscreenCanvasRenderingContext2D; + assert(context); + + context.resetTransform(); + + if (!canvasIsNew) { + if (!alpha && isFirefox()) { + context.fillStyle = 'black'; + context.fillRect(0, 0, width, height); + } else { + context.clearRect(0, 0, width, height); + } + } + + sample.drawWithFit(context, { fit, rotation, crop }); + + const result = { + canvas, + timestamp: sample.timestamp, + duration: sample.duration, + }; + + sample.close(); + return result; + }; +}; + +/** + * An AudioBuffer with additional timing information (timestamp & duration). + * @public + */ +export type WrappedAudioBuffer = { + /** An AudioBuffer. */ + buffer: AudioBuffer; + /** The timestamp of the corresponding audio sample, in seconds. */ + timestamp: number; + /** The duration of the corresponding audio sample, in seconds. */ + duration: number; +}; + +export const audioBufferTransformer = (): SampleTransformer => { + return (sample) => { + const result: WrappedAudioBuffer = { + buffer: sample.toAudioBuffer(), + timestamp: sample.timestamp, + duration: sample.duration, + }; + + sample.close(); + return result; + }; +}; diff --git a/src/input.ts b/src/input.ts index aedf892..7e9ec3a 100644 --- a/src/input.ts +++ b/src/input.ts @@ -6,6 +6,7 @@ * file, You can obtain one at https://mozilla.org/MPL/2.0/. */ +import { SampleCursor } from './cursors'; import { Demuxer } from './demuxer'; import { InputFormat } from './input-format'; import { assert, polyfillSymbolDispose } from './misc'; @@ -44,6 +45,9 @@ export class Input implements Disposable { _reader: Reader; /** @internal */ _disposed = false; + /** @internal */ + // eslint-disable-next-line @typescript-eslint/no-explicit-any + _openSampleCursors = new Set>(); /** True if the input has been disposed. */ get disposed() { @@ -181,7 +185,9 @@ export class Input implements Disposable { this._source._disposed = true; this._source._dispose(); - // TODO this should dispose cursors probably + for (const cursor of [...this._openSampleCursors]) { + void cursor.close(); + } } /** diff --git a/src/media-sink.ts b/src/media-sink.ts index 8d3ab8a..622b39e 100644 --- a/src/media-sink.ts +++ b/src/media-sink.ts @@ -387,7 +387,7 @@ export class EncodedPacketSink { } } -abstract class DecoderWrapper< +export abstract class DecoderWrapper< MediaSample extends VideoSample | AudioSample, > { constructor( @@ -399,6 +399,8 @@ abstract class DecoderWrapper< abstract decode(packet: EncodedPacket): void; abstract flush(): Promise; abstract close(): void; + + abstract get closed(): boolean; } /** @@ -1749,17 +1751,20 @@ export class CanvasSink { } } -class AudioDecoderWrapper extends DecoderWrapper { +export class AudioDecoderWrapper extends DecoderWrapper { decoder: AudioDecoder | null = null; customDecoder: CustomAudioDecoder | null = null; customDecoderCallSerializer = new CallSerializer(); customDecoderQueueSize = 0; + customDecoderClosed = false; // Internal state to accumulate a precise current timestamp based on audio durations, not the (potentially // inaccurate) packet timestamps. currentTimestamp: number | null = null; + onDequeue: (() => unknown) | null = null; + constructor( onSample: (sample: AudioSample) => unknown, onError: (error: Error) => unknown, @@ -1824,6 +1829,10 @@ class AudioDecoderWrapper extends DecoderWrapper { error: onError, }); this.decoder.configure(decoderConfig); + + this.decoder.addEventListener('dequeue', () => { + this.onDequeue?.(); + }); } } @@ -1841,7 +1850,10 @@ class AudioDecoderWrapper extends DecoderWrapper { this.customDecoderQueueSize++; void this.customDecoderCallSerializer .call(() => this.customDecoder!.decode(packet)) - .then(() => this.customDecoderQueueSize--); + .then(() => { + this.customDecoderQueueSize--; + this.onDequeue?.(); + }); } else { assert(this.decoder); this.decoder.decode(packet.toEncodedAudioChunk()); @@ -1859,17 +1871,36 @@ class AudioDecoderWrapper extends DecoderWrapper { close() { if (this.customDecoder) { - void this.customDecoderCallSerializer.call(() => this.customDecoder!.close()); + if (!this.customDecoderClosed) { + this.customDecoderClosed = true; + void this.customDecoderCallSerializer.call(() => this.customDecoder!.close()); + } } else { assert(this.decoder); - this.decoder.close(); + + if (this.decoder.state !== 'closed') { + this.decoder.close(); + } + } + } + + get closed() { + if (this.customDecoder) { + if (this.customDecoderClosed) { + return true; + } + + return !this.customDecoderCallSerializer.errored; + } else { + assert(this.decoder); + return this.decoder.state === 'closed'; } } } // There are a lot of PCM variants not natively supported by the browser and by AudioData. Therefore we need a simple // decoder that maps any input PCM format into a PCM format supported by the browser. -class PcmAudioDecoderWrapper extends DecoderWrapper { +export class PcmAudioDecoderWrapper extends DecoderWrapper { codec: PcmAudioCodec; inputSampleSize: 1 | 2 | 3 | 4 | 8; @@ -1883,6 +1914,10 @@ class PcmAudioDecoderWrapper extends DecoderWrapper { // inaccurate) packet timestamps. currentTimestamp: number | null = null; + isClosed = false; + + onDequeue: (() => unknown) | null = null; + constructor( onSample: (sample: AudioSample) => unknown, onError: (error: Error) => unknown, @@ -2006,6 +2041,8 @@ class PcmAudioDecoderWrapper extends DecoderWrapper { } decode(packet: EncodedPacket) { + this.onDequeue?.(); + const inputView = toDataView(packet.data); const numberOfFrames = packet.byteLength / this.decoderConfig.numberOfChannels / this.inputSampleSize; @@ -2048,7 +2085,11 @@ class PcmAudioDecoderWrapper extends DecoderWrapper { } close() { - // Do nothing + this.isClosed = true; + } + + get closed() { + return this.isClosed; } } diff --git a/src/misc.ts b/src/misc.ts index 2677f8b..d577e13 100644 --- a/src/misc.ts +++ b/src/misc.ts @@ -835,6 +835,8 @@ export const polyfillSymbolDispose = () => { // https://www.typescriptlang.org/docs/handbook/release-notes/typescript-5-2.html // @ts-expect-error Readonly Symbol.dispose ??= Symbol('Symbol.dispose'); + // @ts-expect-error Readonly + Symbol.asyncDispose ??= Symbol('Symbol.asyncDispose'); }; export const isNumber = (x: unknown) => { diff --git a/src/sample.ts b/src/sample.ts index 9d61724..3d0de84 100644 --- a/src/sample.ts +++ b/src/sample.ts @@ -921,6 +921,9 @@ export type AudioSampleCopyToOptions = { * @public */ export class AudioSample implements Disposable { + /** @internal */ + static _openSampleCount = 0; + /** @internal */ _data: AudioData | Uint8Array; /** @internal */ @@ -1033,6 +1036,7 @@ export class AudioSample implements Disposable { this._data = dataBuffer; } + AudioSample._openSampleCount++; finalizationRegistry?.register(this, { type: 'audio', data: this._data }, this); } @@ -1275,6 +1279,7 @@ export class AudioSample implements Disposable { } this._closed = true; + AudioSample._openSampleCount--; } /** diff --git a/test/browser/sample-cursor.test.ts b/test/browser/sample-cursor.test.ts index 60bbfdf..c17e1a2 100644 --- a/test/browser/sample-cursor.test.ts +++ b/test/browser/sample-cursor.test.ts @@ -1,9 +1,15 @@ import { expect, test } from 'vitest'; import { Input } from '../../src/input.js'; -import { UrlSource } from '../../src/source.js'; +import { BufferSource, UrlSource } from '../../src/source.js'; import { ALL_FORMATS } from '../../src/input-format.js'; -import { PacketReader, VideoSampleCursor } from '../../src/cursors.js'; -import { VideoSample } from '../../src/sample.js'; +import { + audioBufferTransformer, + AudioSampleCursor, + canvasTransformer, + PacketReader, + VideoSampleCursor, +} from '../../src/cursors.js'; +import { AudioSample, VideoSample } from '../../src/sample.js'; import { promiseIterateAll } from '../../src/misc.js'; test('Sample cursor seeking', async () => { @@ -14,7 +20,7 @@ test('Sample cursor seeking', async () => { const videoTrack = (await input.getPrimaryVideoTrack())!; const reader = new PacketReader(videoTrack); - const cursor = await VideoSampleCursor.init(reader); + await using cursor = await VideoSampleCursor.init(reader); cursor.debugInfo.enabled = true; expect(cursor.current).toBe(null); @@ -292,19 +298,45 @@ test('Decoder pump error handling', async () => { cursor.debugInfo.enabled = true; cursor.debugInfo.throwInPump = true; - await expect(cursor.seekToFirst()).rejects.toThrow(); + await expect(async () => cursor.seekToFirst()).rejects.toThrow(); + expect(cursor.closed).toBe(true); expect(cursor.current).toBe(null); cursor.debugInfo.throwInPump = false; - const firstPacket = (await cursor.seekToFirst())!; - expect(firstPacket.timestamp).toBe(0); // It has recovered - - await cursor.close(); + await expect(async () => cursor.seekToFirst()).rejects.toThrow(); // It's bricked expect(VideoSample._openSampleCount).toBe(0); }); +test('Decoder errors', 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); + + const cursor1 = await VideoSampleCursor.init(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); + cursor2.debugInfo.enabled = true; + + await cursor2.seekToFirst(); + + cursor2.debugInfo.throwDecoderError = true; + await new Promise(resolve => setTimeout(resolve, 200)); + + expect(() => cursor2.next()).toThrow('Fake decoder error'); + expect(cursor2.closed).toBe(true); +}); + test('Use after close', async () => { using input = new Input({ source: new UrlSource('/trim-buck-bunny.mov'), @@ -313,19 +345,51 @@ test('Use after close', async () => { const videoTrack = (await input.getPrimaryVideoTrack())!; const reader = new PacketReader(videoTrack); - const cursor = await VideoSampleCursor.init(reader); - await cursor.close(); + const cursor1 = await VideoSampleCursor.init(reader); - await expect(async () => await cursor.seekToFirst()).rejects.toThrow(); - await expect(async () => await cursor.seekTo(0)).rejects.toThrow(); - await expect(async () => await cursor.seekToKey(0)).rejects.toThrow(); - await expect(async () => await cursor.next()).rejects.toThrow(); + expect(cursor1.closed).toBe(false); + await cursor1.close(); + expect(cursor1.closed).toBe(true); + + await expect(async () => await cursor1.seekToFirst()).rejects.toThrow('cursor has been closed'); + await expect(async () => await cursor1.seekTo(0)).rejects.toThrow('cursor has been closed'); + 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 commands2 = [ + cursor2.seekToFirst(), + cursor2.next(), + cursor2.close(), + 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'); + + const cursor3 = await VideoSampleCursor.init(reader); + const commands3 = [ + cursor3.seekToFirst(), + cursor3.next(), + cursor3.close(), + 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'); + + expect(VideoSample._openSampleCount).toBe(0); }); test('Command queuing', async () => { using input = new Input({ - source: new UrlSource('/trim-buck-bunny.mov'), + // Fetch the data into RAM to avoid packet lookups causing flaky timing + source: new BufferSource(await fetch('/trim-buck-bunny.mov').then(x => x.arrayBuffer())), formats: ALL_FORMATS, }); @@ -523,9 +587,236 @@ test('Command queuing', async () => { expect(VideoSample._openSampleCount).toBe(0); }); -// TODO: -// - Test throwing if cursor isn't closed at the end of the test -// - Sample cursor mapping function -// - AudioSampleCursor +test('Automatic cursor disposal', async () => { + using input = new Input({ + source: new UrlSource('/trim-buck-bunny.mov'), + formats: ALL_FORMATS, + }); -// Test Sylvie video? + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + + const cursor = await VideoSampleCursor.init(reader); + await cursor.seekToFirst(); + + // No cursor.close() here, but the disposed Input closes the cursor + input.dispose(); + + expect(cursor._closed).toBe(true); +}); + +test('Video with stubborn first sample emit', async () => { + using input = new Input({ + // For some reason, this video is stubborn in the sense that it takes quite a lot of packets for the decoder to + // emit its first samples. This used to cause issues in the past where the decoder got stuck, so good to have it + // tested. + source: new UrlSource('/sylvie-trimmed.mp4'), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + await using cursor = await VideoSampleCursor.init(reader); + + const firstSample = (await cursor.seekToFirst())!; + expect(firstSample).not.toBe(null); + expect(firstSample.timestamp).toBe(0); +}); + +test('AudioSampleCursor', async () => { + using input = new Input({ + source: new UrlSource('/trim-buck-bunny.mov'), + formats: ALL_FORMATS, + }); + + const audioTrack = (await input.getPrimaryAudioTrack())!; + const reader = new PacketReader(audioTrack); + const cursor = await AudioSampleCursor.init(reader); + cursor.debugInfo.enabled = true; + + const firstSample = (await cursor.seekToFirst())!; + expect(firstSample).not.toBe(null); + expect(firstSample.timestamp).toBe(0); + expect(firstSample).toBeInstanceOf(AudioSample); + + const secondSample = (await cursor.next())!; + expect(secondSample.timestamp).toBeCloseTo(firstSample.timestamp + firstSample.duration); + + const thirdSample = (await cursor.seekTo(secondSample.timestamp + 0.05))!; + expect(thirdSample.timestamp).toBeGreaterThan(secondSample.timestamp); + + expect(cursor.debugInfo.pumpsStarted).toBe(1); + + await cursor.seekToFirst(); + + let lastTimestamp = -Infinity; + let total = 0; + for await (const sample of cursor) { + if (total === 0) { + expect(sample.timestamp).toBe(0); + } + + expect(sample.timestamp).toBeGreaterThan(lastTimestamp); + lastTimestamp = sample.timestamp; + total++; + } + + expect(total).toBe(235); + expect(lastTimestamp).toBeCloseTo(5, 1); + + expect(cursor.debugInfo.pumpsStarted).toBe(2); + + const middleSample = (await cursor.seekTo(2.5))!; + expect(middleSample.timestamp).toBeLessThanOrEqual(2.5); + expect(middleSample.timestamp).toBeGreaterThan(2.4); + + expect(cursor.debugInfo.pumpsStarted).toBe(3); + + const commands = [ + cursor.seekTo(0), + cursor.seekTo(0.05), + cursor.seekTo(0.1), + cursor.seekTo(0.15), + cursor.seekTo(0.2), + ]; + // eslint-disable-next-line @typescript-eslint/await-thenable + const result = await Promise.all(commands); + + expect(result[0]!.timestamp).toBeLessThanOrEqual(0); + expect(result[1]!.timestamp).toBeLessThanOrEqual(0.05); + expect(result[2]!.timestamp).toBeLessThanOrEqual(0.1); + expect(result[3]!.timestamp).toBeLessThanOrEqual(0.15); + expect(result[4]!.timestamp).toBeLessThanOrEqual(0.2); + + // One pump was used for all the above commands, even tho they all seek to different key packets + expect(cursor.debugInfo.pumpsStarted).toBe(4); + + await cursor.close(); + + expect(AudioSample._openSampleCount).toBe(0); +}); + +test('Sample mapping', 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); + + let callCount = 0; + await using cursor = await VideoSampleCursor.init(reader, { + transform: (sample) => { + callCount++; + + return { + original: sample, + timestampMs: sample.timestamp * 1000, + }; + }, + }); + + const firstSample = (await cursor.seekToFirst())!; + expect(firstSample.original).toBeInstanceOf(VideoSample); + expect(firstSample.timestampMs).toBe(0); + expect(firstSample.original.closed).toBe(false); + + const secondSample = (await cursor.next())!; + expect(secondSample.original).toBeInstanceOf(VideoSample); + expect(secondSample.timestampMs).toBeCloseTo(firstSample.timestampMs + firstSample.original.duration * 1000); + expect(firstSample.original.closed).toBe(true); + expect(secondSample.original.closed).toBe(false); + + const thirdSample = (await cursor.seekTo(0.5))!; + expect(thirdSample).not.toBe(null); + + expect(callCount).toBe(3); +}); + +test('Canvas transformer', 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); + + const cursor1 = await VideoSampleCursor.init(reader, { + transform: canvasTransformer(), + }); + + const firstSample = (await cursor1.seekToFirst())!; + expect(firstSample.canvas).toBeInstanceOf(HTMLCanvasElement); + expect(firstSample.canvas.width).toBe(videoTrack.displayWidth); + expect(firstSample.canvas.height).toBe(videoTrack.displayHeight); + expect(firstSample.timestamp).toBe(0); + expect(firstSample.duration).toBeGreaterThan(0); + + const nextSample = (await cursor1.next())!; + expect(nextSample.canvas).toBeInstanceOf(HTMLCanvasElement); + expect(nextSample.timestamp).toBeGreaterThan(firstSample.timestamp); + expect(nextSample.canvas).not.toBe(firstSample.canvas); + + await cursor1.close(); + + const cursor2 = await VideoSampleCursor.init(reader, { + transform: canvasTransformer({ + width: 320, + poolSize: 2, + }), + }); + + const sample1 = (await cursor2.seekToFirst())!; + expect(sample1.canvas.width).toBe(320); + expect(sample1.canvas.height).toBe(180); + + const sample2 = (await cursor2.next())!; + expect(sample2.canvas.width).toBe(320); + expect(sample2.canvas.height).toBe(180); + expect(sample2.canvas).not.toBe(sample1.canvas); + + const sample3 = (await cursor2.next())!; + const sample4 = (await cursor2.next())!; + + expect(sample3.canvas).toBe(sample1.canvas); + expect(sample2.canvas).toBe(sample4.canvas); + + await cursor2.close(); + + expect(VideoSample._openSampleCount).toBe(0); +}); + +test('AudioBuffer transformer', async () => { + using input = new Input({ + source: new UrlSource('/trim-buck-bunny.mov'), + formats: ALL_FORMATS, + }); + + const audioTrack = (await input.getPrimaryAudioTrack())!; + const reader = new PacketReader(audioTrack); + + const cursor = await AudioSampleCursor.init(reader, { + transform: audioBufferTransformer(), + }); + + const firstSample = (await cursor.seekToFirst())!; + expect(firstSample.buffer).toBeInstanceOf(AudioBuffer); + expect(firstSample.timestamp).toBe(0); + expect(firstSample.duration).toBeGreaterThan(0); + expect(firstSample.buffer.duration).toBe(firstSample.duration); + + const nextSample = (await cursor.next())!; + expect(nextSample.buffer).toBeInstanceOf(AudioBuffer); + expect(nextSample.timestamp).toBeGreaterThan(firstSample.timestamp); + expect(nextSample.buffer).not.toBe(firstSample.buffer); + expect(nextSample.buffer.duration).toBe(nextSample.duration); + + await cursor.close(); + + expect(AudioSample._openSampleCount).toBe(0); +}); + +// TODO: +// - Then, clean up the cursor.ts code diff --git a/test/node/packet-cursor.test.ts b/test/node/packet-cursor.test.ts index c95aff2..ee49491 100644 --- a/test/node/packet-cursor.test.ts +++ b/test/node/packet-cursor.test.ts @@ -297,7 +297,7 @@ test('Command queuing', async () => { test('verifyKeyPackets with faultily-labeled key frames', async () => { using input = new Input({ - source: new FilePathSource(path.join(__dirname, '../public/fake-cod.mp4')), + source: new FilePathSource(path.join(__dirname, '../public/fucked-keyframes.mp4')), formats: ALL_FORMATS, }); diff --git a/test/public/fake-cod.mp4 b/test/public/fake-cod.mp4 deleted file mode 100644 index 901c608..0000000 Binary files a/test/public/fake-cod.mp4 and /dev/null differ diff --git a/test/public/fucked-keyframes.mp4 b/test/public/fucked-keyframes.mp4 new file mode 100644 index 0000000..5005cae Binary files /dev/null and b/test/public/fucked-keyframes.mp4 differ diff --git a/test/public/sylvie-trimmed.mp4 b/test/public/sylvie-trimmed.mp4 new file mode 100644 index 0000000..18e59ba Binary files /dev/null and b/test/public/sylvie-trimmed.mp4 differ