diff --git a/dev/demux.html b/dev/demux.html index 788a0cf..064ff42 100644 --- a/dev/demux.html +++ b/dev/demux.html @@ -23,12 +23,55 @@ console.log(await cursor.next()); */ + //console.log(await cursor.seekTo(4.9)) + //console.log(await cursor.seekTo(5.1)) + //return; + const mh = [ - cursor.seekTo(2.10), - cursor.seekTo(5), + cursor.seekTo(2.05), + //cursor.seekTo(4.9), + //cursor.seekTo(4.9), + //cursor.seekTo(4.9), + //cursor.seekTo(4.9), + //cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.next(), + cursor.close(), + //cursor.next(), + //cursor.next(), + //cursor.next(), + //cursor.next(), + //cursor.close(), + //cursor.seekTo(2.10), + //cursor.seekTo(5), + //cursor.close(), //cursor.seekTo(2.00), ]; + /* + for (const yo of mh) { + const samp = await yo; + console.log(samp) + samp?.close(); + } + + console.log("done") + */ + console.log(await Promise.all(mh)); //console.log(await cursor.seekTo(2.05)); diff --git a/eslint.config.mjs b/eslint.config.mjs index 29fb45a..aacbf96 100644 --- a/eslint.config.mjs +++ b/eslint.config.mjs @@ -30,6 +30,19 @@ export default tseslint.config( '@typescript-eslint/require-await': 'off', '@stylistic/yield-star-spacing': ['error', { before: false, after: true }], '@typescript-eslint/no-unsafe-enum-comparison': 'off', + '@typescript-eslint/no-unused-vars': [ + 'error', + { + // From https://typescript-eslint.io/rules/no-unused-vars/ + "args": "all", + "argsIgnorePattern": "^_", + "caughtErrors": "all", + "caughtErrorsIgnorePattern": "^_", + "destructuredArrayIgnorePattern": "^_", + "varsIgnorePattern": "^_", + "ignoreRestSiblings": true, + }, + ], }, }, { diff --git a/src/conversion.ts b/src/conversion.ts index 33bbd4d..632b300 100644 --- a/src/conversion.ts +++ b/src/conversion.ts @@ -134,7 +134,8 @@ export type ConversionVideoOptions = { fit?: 'fill' | 'contain' | 'cover'; /** * The angle in degrees to rotate the input video by, clockwise. Rotation is applied before cropping and resizing. - * This rotation is _in addition to_ the natural rotation of the input video as specified in input file's metadata. + * This rotation is _in addition to_ the natural rotation of the input video as specified in the input file's + * metadata. */ rotate?: Rotation; /** @@ -171,7 +172,7 @@ export type ConversionVideoOptions = { * frames improve seeking behavior but increase file size. When using multiple video tracks, you should give them * all the same key frame interval. * - * Setting this fields forces a transcode. + * Setting this field forces a transcode. */ keyFrameInterval?: number; /** When `true`, video will always be re-encoded instead of directly copying over the encoded samples. */ diff --git a/src/cursors.ts b/src/cursors.ts index bedb406..bf57921 100644 --- a/src/cursors.ts +++ b/src/cursors.ts @@ -2,25 +2,76 @@ // Useful in packet context? Or only for sample? import { InputTrack, InputVideoTrack } from './input-track'; -import { PacketRetrievalOptions, VideoDecoderWrapper } from './media-sink'; -import { assert, assertNever, AsyncMutex, AsyncMutex2, insertSorted, last, MaybePromise, promiseWithResolvers, ResultValue, Yo } from './misc'; +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 { EncodedPacket } from './packet'; import { VideoSample } from './sample'; -export class PacketCursor { - track: InputTrack; - _options: PacketRetrievalOptions; - current: EncodedPacket | null = null; - initialized = false; +export class PacketReader { + track: T; - constructor(track: InputTrack, options: PacketRetrievalOptions = {}) { + constructor(track: T) { + if (!(track instanceof InputTrack)) { + throw new TypeError('track must be an InputTrack.'); + } this.track = track; - this._options = options; } - peekAtStart(): MaybePromise { + private maybeVerifyPacketType( + packet: EncodedPacket | null, + options: PacketRetrievalOptions, + ): MaybePromise { + if (!options.verifyKeyPackets || !packet || packet.type === 'delta') { + return packet; + } + + return this.track.determinePacketType(packet).then((determinedType) => { + if (determinedType) { + // @ts-expect-error Technically readonly + packet.type = determinedType; + } + + return packet; + }); + } + + readFirst(options: PacketRetrievalOptions = {}): MaybePromise { + validatePacketRetrievalOptions(options); + const result = new ResultValue(); - const promise = this.track._backing.getFirstPacket(result, this._options); + const promise = this.track._backing.getFirstPacket(result, options); + + if (result.pending) { + return (promise as Promise).then(() => this.maybeVerifyPacketType(result.value, options)); + } else { + return this.maybeVerifyPacketType(result.value, options); + } + } + + readAt(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(() => this.maybeVerifyPacketType(result.value, options)); + } else { + return this.maybeVerifyPacketType(result.value, options); + } + } + + readKeyAt(timestamp: number, options: PacketRetrievalOptions = {}): MaybePromise { + validateTimestamp(timestamp); + validatePacketRetrievalOptions(options); + + if (options.verifyKeyPackets) { + return this.readKeyAtVerified(timestamp, options); + } + + const result = new ResultValue(); + const promise = this.track._backing.getKeyPacket(result, timestamp, options); if (result.pending) { return (promise as Promise).then(() => result.value); @@ -29,24 +80,56 @@ export class PacketCursor { } } - seekToStart(): MaybePromise { + private async readKeyAtVerified( + timestamp: number, + options: PacketRetrievalOptions, + ): Promise { const result = new ResultValue(); - const promise = this.track._backing.getFirstPacket(result, this._options); + const promise = this.track._backing.getKeyPacket(result, timestamp, options); + if (result.pending) await promise; + + const packet = result.value; + if (!packet) { + return null; + } + + const determinedType = await this.track.determinePacketType(packet); + if (determinedType === 'delta') { + // Try returning the previous key packet (in hopes that it's actually a key packet) + return this.readKeyAtVerified(packet.timestamp - 1 / this.track.timeResolution, options); + } + + return packet; + } + + readNext(from: EncodedPacket, options: PacketRetrievalOptions = {}): MaybePromise { + if (!(from instanceof EncodedPacket)) { + throw new TypeError('from must be an EncodedPacket.'); + } + validatePacketRetrievalOptions(options); + + const result = new ResultValue(); + const promise = this.track._backing.getNextPacket(result, from, options); if (result.pending) { - return (promise as Promise).then(() => { - this.initialized = true; - return this.current = result.value; - }); + return (promise as Promise).then(() => this.maybeVerifyPacketType(result.value, options)); } else { - this.initialized = true; - return this.current = result.value; + return this.maybeVerifyPacketType(result.value, options); } } - peekAt(timestamp: number): MaybePromise { + readNextKey(from: EncodedPacket, options: PacketRetrievalOptions = {}): MaybePromise { + if (!(from instanceof EncodedPacket)) { + throw new TypeError('from must be an EncodedPacket.'); + } + validatePacketRetrievalOptions(options); + + if (options.verifyKeyPackets) { + return this.readNextKeyVerified(from, options); + } + const result = new ResultValue(); - const promise = this.track._backing.getPacket(result, timestamp, this._options); + const promise = this.track._backing.getNextKeyPacket(result, from, options); if (result.pending) { return (promise as Promise).then(() => result.value); @@ -55,144 +138,498 @@ export class PacketCursor { } } + private async readNextKeyVerified( + from: EncodedPacket, + options: PacketRetrievalOptions, + ): Promise { + const result = new ResultValue(); + const promise = this.track._backing.getNextKeyPacket(result, from, options); + if (result.pending) await promise; + + const nextPacket = result.value; + if (!nextPacket) { + return null; + } + + const determinedType = await this.track.determinePacketType(nextPacket); + if (determinedType === 'delta') { + // Try returning the next key packet (in hopes that it's actually a key packet) + return this.readNextKeyVerified(nextPacket, options); + } + + return nextPacket; + } +} + +export class PacketCursor { + reader: PacketReader; + options: PacketRetrievalOptions; + + current: EncodedPacket | null = null; + nextIsFirst = true; + callSerializer = new CallSerializer2(); + + constructor(reader: PacketReader, options: PacketRetrievalOptions = {}) { + if (!(reader instanceof PacketReader)) { + throw new TypeError('reader must be a PacketReader.'); + } + validatePacketRetrievalOptions(options); + + this.reader = reader; + this.options = options; + } + + private seekToFirstDirect(): MaybePromise { + const result = this.reader.readFirst(this.options); + + const onPacket = (packet: EncodedPacket | null) => { + this.nextIsFirst = false; + return this.current = packet; + }; + + if (result instanceof Promise) { + return result.then(onPacket); + } else { + return onPacket(result); + } + } + + seekToFirst(): MaybePromise { + return this.callSerializer.call(() => this.seekToFirstDirect()); + } + seekTo(timestamp: number): MaybePromise { - const result = new ResultValue(); - const promise = this.track._backing.getPacket(result, timestamp, this._options); + return this.callSerializer.call(() => { + const result = this.reader.readAt(timestamp, this.options); - if (result.pending) { - return (promise as Promise).then(() => { - this.initialized = true; - return this.current = result.value; - }); - } else { - this.initialized = true; - return this.current = result.value; - } - } + const onPacket = (packet: EncodedPacket | null) => { + this.nextIsFirst = !packet; + return this.current = packet; + }; - peekKeyAt(timestamp: number): MaybePromise { - const result = new ResultValue(); - const promise = this.track._backing.getKeyPacket(result, timestamp, this._options); - - if (result.pending) { - return (promise as Promise).then(() => result.value); - } else { - return result.value; - } + if (result instanceof Promise) { + return result.then(onPacket); + } else { + return onPacket(result); + } + }); } seekToKey(timestamp: number): MaybePromise { - const result = new ResultValue(); - const promise = this.track._backing.getKeyPacket(result, timestamp, this._options); + return this.callSerializer.call(() => { + const result = this.reader.readKeyAt(timestamp, this.options); - if (result.pending) { - return (promise as Promise).then(() => { - this.initialized = true; - return this.current = result.value; - }); - } else { - this.initialized = true; - return this.current = result.value; - } - } + const onPacket = (packet: EncodedPacket | null) => { + this.nextIsFirst = !packet; + return this.current = packet; + }; - _ensureInitialized() { - if (!this.initialized) { - throw new Error('You must first initialize the cursor to a position by calling any of the seek methods.'); - } + if (result instanceof Promise) { + return result.then(onPacket); + } else { + return onPacket(result); + } + }); } next(): MaybePromise { - this._ensureInitialized(); + return this.callSerializer.call(() => { + if (this.nextIsFirst) { + return this.seekToFirstDirect(); + } - if (!this.current) { - return null; - } + if (!this.current) { + return null; + } - const result = new ResultValue(); - const promise = this.track._backing.getNextPacket(result, this.current, this._options); + const result = this.reader.readNext(this.current, this.options); - if (result.pending) { - return (promise as Promise).then(() => { - return this.current = result.value; - }); - } else { - return this.current = result.value; - } - } + const onPacket = (packet: EncodedPacket | null) => { + return this.current = packet; + }; - peekNextKey(): MaybePromise { - this._ensureInitialized(); - - if (!this.current) { - return null; - } - - const result = new ResultValue(); - const promise = this.track._backing.getNextKeyPacket(result, this.current, this._options); - - if (result.pending) { - return (promise as Promise).then(() => result.value); - } else { - return result.value; - } + if (result instanceof Promise) { + return result.then(onPacket); + } else { + return onPacket(result); + } + }); } nextKey(): MaybePromise { - this._ensureInitialized(); + return this.callSerializer.call(() => { + if (this.nextIsFirst) { + return this.seekToFirstDirect(); + } - if (!this.current) { - return null; - } + if (!this.current) { + return null; + } - const result = new ResultValue(); - const promise = this.track._backing.getNextKeyPacket(result, this.current, this._options); + const result = this.reader.readNextKey(this.current, this.options); - if (result.pending) { - return (promise as Promise).then(() => { - return this.current = result.value; - }); - } else { - return this.current = result.value; - } + const onPacket = (packet: EncodedPacket | null) => { + return this.current = packet; + }; + + if (result instanceof Promise) { + return result.then(onPacket); + } else { + return onPacket(result); + } + }); } - async iterate(callback: (packet: EncodedPacket, stop: () => void) => MaybePromise) { - this._ensureInitialized(); - + async iterate( + callback: (packet: EncodedPacket, stop: () => void) => MaybePromise, + ) { let stopped = false; const stop = () => stopped = true; - while (this.current) { - const result = callback(this.current, stop); - if (result instanceof Promise) await result; + 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; } - let next = this.next(); - if (next instanceof Promise) next = await next; + const result = this.next(); + if (result instanceof Promise) await result; - this.current = next; + if (!this.current) { + break; + } } } // eslint-disable-next-line @stylistic/generator-star-spacing async *[Symbol.asyncIterator]() { - this._ensureInitialized(); + const donePromise = this.callSerializer.done(); + if (donePromise) await donePromise; - while (this.current) { - yield this.current; + while (true) { + if (this.current) { + yield this.current; + } - let next = this.next(); - if (next instanceof Promise) next = await next; + const result = this.next(); + if (result instanceof Promise) await result; - this.current = next; + if (!this.current) { + break; + } } } + + waitUntilIdle() { + return this.callSerializer.done(); + } } +/* +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); @@ -200,50 +637,70 @@ for (const timestamp of timestamps) { } */ -// Problem: Decoding further than is needed. This is not a problem for playback-like operations where it is desired to -// maintain a buffer of pre-decoded frames. It breaks down when using the cursor for quick, pinpointed seek operations. -// It would be cool if the cursor could still be used for pinpointed seeks. This provides a unified API for both -// sequential as well as random access patterns. -// Idea: By default, queue only the packets that are needed to reach a packet. Then, when encode size reaches zero and -// the thing has not been yielded yet, give it another group of packets. -// But that doesn't properly fill the decoder queue when just iterating using next(). -// Idea: Calling next() "unlocks" the decoder, keeping a healthy queue into the future until seek is called? I don't -// know, feels weird. What I like better is something like, queue until the target packet first, then when the -// decode queue drops to zero and there are no new requests, queue more packets. This means that in moments of silence, -// the next samples will be queued. I think this might be better than going off of next(). This is because I can think -// of use cases that don't need next() but still perform playback-like operations using ONLY seeks. -// Or: make it configurable? Like some sort of queueLength or maxQueueLength? Then it's controllable but also adds -// another layer the user has to "know" to do. +type PendingRequest = { + timestamp: number; + promise: Promise; + resolve: (sample: VideoSample | null) => void; + reject: (error: unknown) => void; + successor: PendingRequest | null; +}; -export class VideoSampleCursor2 { - track: InputVideoTrack; - initialized = false; +type VideoSampleCursorOptions = { + autoClose?: boolean; +}; + +export class VideoSampleCursor { + packetReader: PacketReader; packetCursor: PacketCursor; + options: VideoSampleCursorOptions; + autoClose: boolean; pumpRunning = false; decoder: VideoDecoderWrapper; - // current: VideoSample | null = null; + current: VideoSample | null = null; sampleQueue: VideoSample[] = []; queueDequeue = promiseWithResolvers(); - pendingRequests: { - timestamp: number; - resolve: (sample: VideoSample | null) => void; - }[] = []; + pendingRequests: PendingRequest[] = []; + lastPendingRequest: PendingRequest | null = null; - stopPump = false; + predictedRequests = 0; + nextIsFirst = true; + + pumpStopQueued = false; pumpStopped = promiseWithResolvers(); decodedTimestamps: number[] = []; maxDecodedSequenceNumber = -1; - pumpTargetSequenceNumber = -1; - pumpMutex = new AsyncMutex2(); + pumpTarget: EncodedPacket | null = null; + pumpMutex = new AsyncMutex3(); + closed = false; - private constructor(track: InputVideoTrack, decoder: VideoDecoderWrapper) { - this.track = track; + debugInfo = { + enabled: false, + pumpsStarted: 0, + seekPackets: [] as (EncodedPacket | null)[], + decodedPackets: [] as EncodedPacket[], + throwInPump: false, + }; + + private constructor( + reader: PacketReader, + decoder: VideoDecoderWrapper, + options: VideoSampleCursorOptions, + ) { + this.packetReader = reader; this.decoder = decoder; - this.packetCursor = new PacketCursor(track); + this.packetCursor = new PacketCursor(reader); + this.options = options; + this.autoClose = options.autoClose ?? true; } - static async init(track: InputVideoTrack) { + 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.'); + } + + 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' @@ -261,10 +718,8 @@ export class VideoSampleCursor2 { cursor.decodedTimestamps.shift(); } - console.log('revc', sample.timestamp, cursor.pendingRequests.length); - if (cursor.pendingRequests.length === 0) { - if (cursor.stopPump) { + if (cursor.pumpStopQueued) { sample.close(); } else { cursor.sampleQueue.push(sample); @@ -274,10 +729,18 @@ export class VideoSampleCursor2 { for (let i = 0; i < cursor.pendingRequests.length; i++) { const request = cursor.pendingRequests[i]!; - if (request.timestamp <= sample.timestamp) { - request.resolve(given ? sample.clone() : sample); - cursor.pendingRequests.splice(i--, 1); - given = true; + 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++; } } @@ -286,13 +749,11 @@ export class VideoSampleCursor2 { } } - // cursor.current?.close(); - // cursor.current = sample; - cursor.queueDequeue.resolve(); cursor.queueDequeue = promiseWithResolvers(); }, (error) => { + // TODO THIS TODO THIS TODO THIS TODO THIS console.error(error); }, track.codec, @@ -301,116 +762,202 @@ export class VideoSampleCursor2 { track.timeResolution, ); - const cursor = new VideoSampleCursor2(track, decoder); + decoder.onDequeue = () => { + cursor.queueDequeue.resolve(); + cursor.queueDequeue = promiseWithResolvers(); + }; + + const cursor = new VideoSampleCursor(reader, decoder, options); return cursor; } + setCurrent(newCurrent: VideoSample | null) { + if (this.autoClose && this.current && this.current !== newCurrent) { + this.current.close(); + } + + this.current = newCurrent; + } + getNextExpectedTimestamp() { if (this.sampleQueue.length > 0) { return this.sampleQueue[0]!.timestamp; } } - async seekTo(timestamp: number): Promise { - this.initialized = true; // too late? + _ensureNotClosed() { + if (this.closed) { + throw new Error('This cursor has been closed and can no longer be used.'); + } + } - console.log('a'); - while (this.pumpMutex.locked) { - console.log('waiting...'); - await this.pumpMutex.promise; + async _seekToPacket( + res: ResultValue, + targetPacketPromise: MaybePromise, + lock?: AsyncMutexLock, + ): Promise { + this.predictedRequests++; + + if (!lock) { + const mutexPromise = this.pumpMutex.request(); + if (mutexPromise) { + await mutexPromise; + } + lock = this.pumpMutex.lock(); } - console.log('GOIN IN'); + using deferred = defer(() => { + lock?.release(); + this.predictedRequests--; + }); - using _ = this.pumpMutex.lock(); + const targetPacket = targetPacketPromise instanceof Promise + ? await targetPacketPromise + : targetPacketPromise; + + if (this.debugInfo.enabled) { + this.debugInfo.seekPackets.push(targetPacket); + } + + this.nextIsFirst = !targetPacket; - const targetPacket = await this.packetCursor.peekAt(timestamp); - console.log('target'); if (!targetPacket) { - return null; + this.setCurrent(null); + return res.set(null); + } + + if (this.current?.timestamp === targetPacket.timestamp) { + return res.set(this.current); } let setNewPump = true; - if (this.sampleQueue.length > 0 && targetPacket.timestamp <= this.sampleQueue[0]!.timestamp) { - console.log('This bitch case kicked'); + if (this.sampleQueue.length > 0 && targetPacket.timestamp < this.sampleQueue[0]!.timestamp) { + } else { while (this.sampleQueue.length > 0) { const nextSample = this.sampleQueue[0]!; if (targetPacket.timestamp <= nextSample.timestamp) { - console.log('used this path'); - return nextSample; + this.setCurrent(nextSample); + return res.set(nextSample); } this.sampleQueue.shift(); this.queueDequeue.resolve(); this.queueDequeue = promiseWithResolvers(); + nextSample.close(); } - if (this.maxDecodedSequenceNumber !== -1) { - // This means a packet was queued for decode and the cursor is initialized + if (this.pumpTarget) { + const max = Math.max(this.pumpTarget.sequenceNumber, this.maxDecodedSequenceNumber); - if (targetPacket.sequenceNumber <= this.maxDecodedSequenceNumber) { - const nextExpectedTimestamp = this.decodedTimestamps[0]; - if (!nextExpectedTimestamp || nextExpectedTimestamp > timestamp) { + if (targetPacket.sequenceNumber <= max) { + const nextExpectedTimestamp = this.decodedTimestamps[0] ?? this.packetCursor.current?.timestamp; + if (nextExpectedTimestamp === undefined || nextExpectedTimestamp > targetPacket.timestamp) { // yeah + // formulate the "no nextExpectedTimestamp" case } else { setNewPump = false; } } else { - const key = await this.packetCursor.peekNextKey(); - if (!key || targetPacket.sequenceNumber < key.sequenceNumber) { + if (targetPacket.timestamp - this.pumpTarget.timestamp < 0.1) { setNewPump = false; + } else { + let key = this.packetReader.readNextKey(this.pumpTarget, { verifyKeyPackets: true }); + if (key instanceof Promise) key = await key; + + if ( + !key + || targetPacket.sequenceNumber < key.sequenceNumber + ) { + setNewPump = false; + } } } } } + if (setNewPump && this.pumpRunning) { + await this.stopPump(); + } + + if (!this.pumpTarget || targetPacket.sequenceNumber > this.pumpTarget.sequenceNumber) { + this.pumpTarget = targetPacket; + } + if (setNewPump) { - console.log('setting up a new PUMP'); + const result = this.packetCursor.seekToKey(targetPacket.timestamp); + if (result instanceof Promise) await result; - if (this.pumpRunning) { - this.stopPump = true; - this.queueDequeue.resolve(); - this.queueDequeue = promiseWithResolvers(); - await this.pumpStopped.promise; - - for (const sample of this.sampleQueue) { - sample.close(); - } - this.sampleQueue.length = 0; - this.maxDecodedSequenceNumber = -1; - this.decodedTimestamps.length = 0; - this.pumpTargetSequenceNumber = -1; - this.stopPump = false; - } - - await this.packetCursor.seekToKey(timestamp); void this.runPump(); - await Promise.resolve(); // lol } const request = promiseWithResolvers(); - this.pendingRequests.push({ + const pendingRequest: PendingRequest = { timestamp: targetPacket.timestamp, + promise: request.promise, resolve: request.resolve, - }); - this.pumpTargetSequenceNumber = Math.max(this.pumpTargetSequenceNumber, targetPacket.sequenceNumber); + reject: request.reject, + successor: null, + }; + this.pendingRequests.push(pendingRequest); + this.pendingRequests.sort((a, b) => a.timestamp - b.timestamp); + this.lastPendingRequest = pendingRequest; - console.log('requesting', targetPacket.timestamp); - console.log('Done', this.pumpTargetSequenceNumber); + this.queueDequeue.resolve(); + this.queueDequeue = promiseWithResolvers(); - return request.promise; + deferred.execute(); // Waiting for the return would be too long + + return res.set(await request.promise); } - async next() { - while (this.pumpMutex.locked) { - console.log('waiting next...'); - await this.pumpMutex.promise; - } + seekToFirst(): MaybePromise { + this._ensureNotClosed(); - if (!this.initialized) { - throw new Error('This shud be the indicator the next not being available I think'); + const result = new ResultValue(); + const promise = this._seekToPacket(result, this.packetReader.readFirst()); + + if (result.pending) { + return promise.then(() => result.value); + } else { + return result.value; + } + } + + seekTo(timestamp: number): MaybePromise { + this._ensureNotClosed(); + + 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; + } + } + + seekToKey(timestamp: number): MaybePromise { + this._ensureNotClosed(); + + 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; + } + } + + async _nextInternal(res: ResultValue): Promise { + const mutexPromise = this.pumpMutex.request(); + if (mutexPromise) await mutexPromise; + using lock = this.pumpMutex.lock(); + + if (this.nextIsFirst) { + return await this._seekToPacket(res, this.packetReader.readFirst(), lock); } if (this.sampleQueue.length > 0) { @@ -418,236 +965,247 @@ export class VideoSampleCursor2 { this.queueDequeue.resolve(); this.queueDequeue = promiseWithResolvers(); - return nextSample; + this.setCurrent(nextSample); + return res.set(nextSample); } if (!this.pumpRunning) { - return null; // None more after this, boy + this.setCurrent(null); + return res.set(null); // None more after this, boy } + assert(this.lastPendingRequest); + const request = promiseWithResolvers(); - this.pendingRequests.push({ - timestamp: -Infinity, // Matches any sample timestamp, so any next one will match + const pendingRequest: PendingRequest = { + timestamp: -Infinity, + promise: request.promise, resolve: request.resolve, - }); + reject: request.reject, + successor: null, + }; - return request.promise; - } + if (this.pendingRequests.length === 0) { + this.pendingRequests.push(pendingRequest); + } else { + this.lastPendingRequest.successor = pendingRequest; + } + this.lastPendingRequest = pendingRequest; - async runPump() { - assert(this.packetCursor.current); - - this.pumpRunning = true; - - while ( - this.packetCursor.current - && (!this.stopPump || this.packetCursor.current.sequenceNumber <= this.pumpTargetSequenceNumber) - ) { - const maxQueueSize = 8 ?? computeMaxQueueSize(this.sampleQueue.length); // temp - if (!this.stopPump && this.sampleQueue.length + this.decoder.getDecodeQueueSize() > maxQueueSize) { - await this.queueDequeue.promise; - continue; - } - - console.log('DECODING', this.packetCursor.current.timestamp); - - insertSorted(this.decodedTimestamps, this.packetCursor.current.timestamp, x => x); - this.maxDecodedSequenceNumber = this.packetCursor.current.sequenceNumber; - this.decoder.decode(this.packetCursor.current); - await this.packetCursor.next(); + // Note that the next packet we get here is not necessarily the packet belonging to the next sample, since we + // can have out of order timestamps when B-frames are at play. However, if next() is called a sufficiently + // large amount of times, then the pump target will stay roughly in sync with the desired next sample. + const next = this.pumpTarget && await this.packetReader.readNext(this.pumpTarget, { metadataOnly: true }); + if (next) { + this.pumpTarget = next; } - console.log('stopping current pump...'); - await this.decoder.flush(); + lock.release(); // Waiting for the return would be too long - this.pumpStopped.resolve(); - this.pumpStopped = promiseWithResolvers(); - - this.pendingRequests.forEach(x => x.resolve(null)); - this.pendingRequests.length = 0; - - this.pumpRunning = false; - console.log('pump stopped/ended'); - } -} - -async function weJustTesting() { - const cursor = new VideoSampleCursor2(); - - // Spins up decoder and resolves to the sample - await cursor.seekTo(2); - - // This can do multiple things: - // - It pops its internal sample queue until it finds a matching frame; in this case, it returns instantly (no promise) - // - If that wasn't possible, but the packet that corresponds to the requested sample was already queued for encoding, - // it will wait for the decoder to spit it out and then returns it. I guess this requires a "pending requests" ahh - // structure somewhere. - // - If that's also not the case, but the seeked packet is in the current GOP, then it just keeps pumping packets into - // the decoder. - // - If the requested packet is outside of the current GOP or "backwards" from the current stream, it resets the - // internal decoder - // In any case, there's always a "pump" running that supplies the decoder with new packets to decode. This pump is - // halted if the internal queue is sufficiently large, and is resumed when samples are consumed. - // This pump is reset if necessary. - await cursor.seekTo(2.1); -} - -export class VideoSampleCursor { - track: InputVideoTrack; - current: EncodedPacket | null = null; - initialized = false; - packetCursor: PacketCursor; - packetCursor2: PacketCursor; - - decoder!: VideoDecoderWrapper; - - constructor(track: InputVideoTrack) { - this.track = track; - this.packetCursor = new PacketCursor(track, { verifyKeyPackets: true }); - this.packetCursor2 = new PacketCursor(track, { verifyKeyPackets: true }); // not good + return res.set(await request.promise); } - async init() { - if (!(await this.track.canDecode())) { - throw new Error( - 'This video track cannot be decoded by this browser. Make sure to check decodability before using' - + ' a track.', - ); + next(): MaybePromise { + this._ensureNotClosed(); + + const result = new ResultValue(); + const promise = this._nextInternal(result); + + if (result.pending) { + return promise.then(() => result.value); + } else { + return result.value; } - - const decoderConfig = await this.track.getDecoderConfig(); - - this.decoder = new VideoDecoderWrapper( - (sample) => { - this.sampleQueue.push(sample); - }, - (error) => { - // Un que? - }, - this.track.codec!, - decoderConfig!, - this.track.rotation, - this.track.timeResolution, - ); } - pumpFinished = promiseWithResolvers(); - queueDequeue = promiseWithResolvers(); - terminatePump = false; - pumpRunning = false; - sampleQueue: VideoSample[] = []; + async iterate( + callback: (packet: VideoSample, stop: () => void) => MaybePromise, + ) { + let stopped = false; + const stop = () => stopped = true; - async runPump() { - this.pumpRunning = true; + const waitPromise = this.waitUntilIdle(); + if (waitPromise) await waitPromise; - while (this.packetCursor.current && !this.terminatePump) { - const maxQueueSize = computeMaxQueueSize(0); - if (0 + this.decoder.getDecodeQueueSize() > maxQueueSize) { - this.queueDequeue = promiseWithResolvers(); - await this.queueDequeue.promise; - continue; + while (true) { + if (this.current) { + const result = callback(this.current, stop); + if (result instanceof Promise) await result; } - this.decoder.decode(this.packetCursor.current); + if (stopped) { + break; + } - const result = this.packetCursor.next(); + const result = this.next(); if (result instanceof Promise) await result; - } - await this.decoder.flush(); - this.pumpFinished.resolve(); - this.pumpRunning = false; + if (!this.current) { + break; + } + } } - async beginNewRun() { - if (this.pumpRunning) { - this.terminatePump = true; - await this.pumpFinished.promise; + // eslint-disable-next-line @stylistic/generator-star-spacing + async *[Symbol.asyncIterator]() { + const waitPromise = this.waitUntilIdle(); + if (waitPromise) await waitPromise; + while (true) { + if (this.current) { + yield this.current; + } + + const result = this.next(); + if (result instanceof Promise) await result; + + if (!this.current) { + break; + } + } + } + + waitUntilIdle() { + if (this.pendingRequests.length === 0) { + return null; + } + + let lastRequest = last(this.pendingRequests)!; + while (lastRequest.successor) { + lastRequest = lastRequest.successor; + } + + return lastRequest.promise + .catch(() => {}) + .then(() => {}); + } + + closePromise: Promise | null = null; + + close() { + return this.closePromise ??= (async () => { + this.predictedRequests++; + + const mutexPromise = this.pumpMutex.request(); + if (mutexPromise) await mutexPromise; + + this.closed = true; + + if (this.pumpRunning) { + await this.stopPump(); + } + + this.setCurrent(null); + this.decoder.close(); + })(); + } + + async stopPump() { + assert(this.pumpRunning); + + this.pumpStopQueued = true; + this.queueDequeue.resolve(); + this.queueDequeue = promiseWithResolvers(); + await this.pumpStopped.promise; + } + + async runPump() { + try { + assert(this.packetCursor.current); + assert(this.pumpTarget); + + this.pumpRunning = true; + + if (this.debugInfo.enabled) { + this.debugInfo.pumpsStarted++; + } + + while ( + this.packetCursor.current + && ( + !this.pumpStopQueued + || this.packetCursor.current.sequenceNumber <= this.pumpTarget.sequenceNumber + || this.pendingRequests.some(x => x.successor || x.timestamp === -Infinity) + ) + ) { + if (this.debugInfo.enabled && this.debugInfo.throwInPump) { + throw new Error('Throwing artificially.'); + } + + if ( + this.packetCursor.current.sequenceNumber > this.pumpTarget.sequenceNumber + && this.predictedRequests > 0 + && !this.pendingRequests.some(x => x.successor || x.timestamp === -Infinity) + ) { + await this.queueDequeue.promise; + continue; + } + + const decodeQueueSize = this.decoder.getDecodeQueueSize(); + if (this.sampleQueue.length + decodeQueueSize >= 4) { + await this.queueDequeue.promise; + continue; + } + + insertSorted(this.decodedTimestamps, this.packetCursor.current.timestamp, x => x); + this.maxDecodedSequenceNumber = this.packetCursor.current.sequenceNumber; + this.decoder.decode(this.packetCursor.current); + + if (this.debugInfo.enabled) { + this.debugInfo.decodedPackets.push(this.packetCursor.current); + } + + await this.packetCursor.next(); + } + + if (this.pendingRequests.length > 0 || !this.closed) { + await this.decoder.flush(); + } + + const uh = (request: PendingRequest) => { + request.resolve(null); + if (request.successor) { + uh(request.successor); + } + }; + + this.pendingRequests.forEach(uh); + this.setCurrent(null); + } catch (error) { + if (!this.decoder.closed) { + if (this.pendingRequests.length > 0 || !this.closed) { + 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 + } + } finally { for (const sample of this.sampleQueue) { sample.close(); } this.sampleQueue.length = 0; + + this.pendingRequests.length = 0; + this.lastPendingRequest = null; + this.pumpStopped.resolve(); + this.pumpStopped = promiseWithResolvers(); + this.pumpRunning = false; + this.maxDecodedSequenceNumber = -1; + this.decodedTimestamps.length = 0; + this.pumpTarget = null; + this.pumpStopQueued = false; } - - // todo errors - void this.runPump(); - } - - async _seekToCurrentPacket(res: ResultValue): Promise { - const targetPacket = this.packetCursor.current; - assert(targetPacket); - - if (targetPacket.type !== 'key') { - await this.packetCursor.seekToKey(targetPacket.timestamp); - } - } - - /* - seekToStart(): MaybePromise { - const onPacket = (packet: EncodedPacket | null) => { - if (!packet) { - return null; - } - - const result = new ResultValue(); - const promise = this._seekToPacket(result, packet); - - if (result.pending) { - return promise.then(() => result.value); - } else { - return result.value; - } - }; - - const packet = this.packetCursor.seekToStart(); - if (packet instanceof Promise) { - return packet.then(onPacket); - } else { - return onPacket(packet); - } - } - */ - - async seekTo(timestamp: number): Promise { - const packet = await this.packetCursor2.seekTo(timestamp); - if (!packet) { - return null; // I guess? - } - - if (this.packetCursor.current) { - // if (packet.sequenceNumber) - } - - /* - const onPacket = (packet: EncodedPacket | null) => { - if (!packet) { - return null; - } - - const result = new ResultValue(); - const promise = this._seekToPacket(result, packet); - - if (result.pending) { - return promise.then(() => result.value); - } else { - return result.value; - } - }; - - const packet = this.packetCursor.seekTo(timestamp); - if (packet instanceof Promise) { - return packet.then(onPacket); - } else { - return onPacket(packet); - } - */ } } - -const computeMaxQueueSize = (decodedSampleQueueSize: number) => { - // If we have decoded samples lying around, limit the total queue size to a small value (decoded samples can use up - // a lot of memory). If not, we're fine with a much bigger queue of encoded packets waiting to be decoded. In fact, - // some decoders only start flushing out decoded chunks when the packet queue is large enough. - return decodedSampleQueueSize === 0 ? 40 : 8; -}; diff --git a/src/encode.ts b/src/encode.ts index 9f07a69..006bf0d 100644 --- a/src/encode.ts +++ b/src/encode.ts @@ -103,7 +103,7 @@ export const validateVideoEncodingConfig = (config: VideoEncodingConfig) => { }; /** - * Additional options that control audio encoding. + * Additional options that control video encoding. * @group Encoding * @public */ diff --git a/src/index.ts b/src/index.ts index ec6454a..acd42e6 100644 --- a/src/index.ts +++ b/src/index.ts @@ -184,7 +184,7 @@ export { } from './media-sink'; export { PacketCursor, - VideoSampleCursor2, + VideoSampleCursor, } from './cursors'; export { Conversion, diff --git a/src/input.ts b/src/input.ts index 5e78243..aedf892 100644 --- a/src/input.ts +++ b/src/input.ts @@ -180,6 +180,8 @@ export class Input implements Disposable { this._source._disposed = true; this._source._dispose(); + + // TODO this should dispose cursors probably } /** diff --git a/src/media-sink.ts b/src/media-sink.ts index db36e60..c4bf404 100644 --- a/src/media-sink.ts +++ b/src/media-sink.ts @@ -67,7 +67,7 @@ export type PacketRetrievalOptions = { verifyKeyPackets?: boolean; }; -const validatePacketRetrievalOptions = (options: PacketRetrievalOptions) => { +export const validatePacketRetrievalOptions = (options: PacketRetrievalOptions) => { if (!options || typeof options !== 'object') { throw new TypeError('options must be an object.'); } @@ -82,7 +82,7 @@ const validatePacketRetrievalOptions = (options: PacketRetrievalOptions) => { } }; -const validateTimestamp = (timestamp: number) => { +export const validateTimestamp = (timestamp: number) => { if (!isNumber(timestamp)) { throw new TypeError('timestamp must be a number.'); // It can be non-finite, that's fine } @@ -830,6 +830,7 @@ export class VideoDecoderWrapper extends DecoderWrapper { customDecoder: CustomVideoDecoder | null = null; customDecoderCallSerializer = new CallSerializer(); customDecoderQueueSize = 0; + customDecoderClosed = false; inputTimestamps: number[] = []; // Timestamps input into the decoder, sorted. sampleQueue: VideoSample[] = []; // Safari-specific thing, check usage. @@ -850,6 +851,8 @@ export class VideoDecoderWrapper extends DecoderWrapper { currentAlphaPacketIndex = 0; alphaRaslSkipped = false; // For HEVC stuff + onDequeue: (() => unknown) | null = null; + constructor( onSample: (sample: VideoSample) => unknown, onError: (error: Error) => unknown, @@ -918,6 +921,10 @@ export class VideoDecoderWrapper extends DecoderWrapper { error: onError, }); this.decoder.configure(this.decoderConfig); + + this.decoder.addEventListener('dequeue', () => { + this.onDequeue?.(); + }); } } @@ -949,7 +956,10 @@ export class VideoDecoderWrapper extends DecoderWrapper { this.customDecoderQueueSize++; void this.customDecoderCallSerializer .call(() => this.customDecoder!.decode(packet)) - .then(() => this.customDecoderQueueSize--); + .then(() => { + this.customDecoderQueueSize--; + this.onDequeue?.(); + }); } else { assert(this.decoder); @@ -1026,6 +1036,10 @@ export class VideoDecoderWrapper extends DecoderWrapper { error: this.onError, }); this.alphaDecoder.configure(this.decoderConfig); + + this.alphaDecoder.addEventListener('dequeue', () => { + this.onDequeue?.(); + }); } const type = determineVideoPacketType(this.codec, this.decoderConfig, packet.sideData.alpha); @@ -1186,11 +1200,19 @@ export class VideoDecoderWrapper 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(); - this.alphaDecoder?.close(); + + if (this.decoder.state !== 'closed') { + this.decoder.close(); + } + if (this.alphaDecoder && this.alphaDecoder.state !== 'closed') { + this.alphaDecoder.close(); + } this.colorQueue.forEach(x => x.close()); this.colorQueue.length = 0; @@ -1205,6 +1227,19 @@ export class VideoDecoderWrapper extends DecoderWrapper { } this.sampleQueue.length = 0; } + + get closed() { + if (this.customDecoder) { + if (this.customDecoderClosed) { + return true; + } + + return !this.customDecoderCallSerializer.errored; + } else { + assert(this.decoder); + return this.decoder.state === 'closed'; + } + } } /** Utility class that merges together color and alpha information using simple WebGL 2 shaders. */ diff --git a/src/metadata.ts b/src/metadata.ts index dc1d166..869e273 100644 --- a/src/metadata.ts +++ b/src/metadata.ts @@ -72,7 +72,7 @@ export type MetadataTags = { * - Ogg: The key-value string pairs from the Vorbis-style comment header (see RFC 7845, Section 5.2). * Additionally, the `'vendor'` key refers to the vendor string within this header. * - WAVE: The individual metadata chunks within the RIFF INFO chunk. Values are always ISO 8859-1 strings. - * - FLAC: The key-value string pairs from the vorbis metadata block (see RFC 9639, Section D.2.3). + * - FLAC: The key-value string pairs from the Vorbis metadata block (see RFC 9639, Section D.2.3). * Additionally, the `'vendor'` key refers to the vendor string within this header. */ raw?: Record; diff --git a/src/misc.ts b/src/misc.ts index cad030e..29c04c2 100644 --- a/src/misc.ts +++ b/src/misc.ts @@ -666,9 +666,15 @@ export const computeRationalApproximation = (x: number, maxDenominator: number) export class CallSerializer { currentPromise = Promise.resolve(); + errored = false; call(fn: () => Promise | void) { - return this.currentPromise = this.currentPromise.then(fn); + return this.currentPromise = this.currentPromise + .then(fn) + .catch((error) => { + this.errored = true; + throw error; + }); } } @@ -840,11 +846,143 @@ export class AsyncMutex2 { const { promise, resolve } = promiseWithResolvers(); this.promise = promise; + let released = false; + return { - [Symbol.dispose]: () => { + release: () => { + if (released) { + return; + } + released = true; + resolve(); this.locked = false; }, + [Symbol.dispose]() { + this.release(); + }, }; } } + +export type AsyncMutexLock = { + release: () => void; + [Symbol.dispose]: () => void; +}; + +export class AsyncMutex3 { + locked = false; + resolverQueue: (() => void)[] = []; + + lock(): AsyncMutexLock { + if (this.locked) { + throw new Error('Mutex already locked.'); + } + + this.locked = true; + let released = false; + + return { + release: () => { + if (released) { + return; + } + released = true; + + this.locked = false; + + if (this.resolverQueue.length > 0) { + const resolve = this.resolverQueue.shift()!; + resolve(); + } + }, + [Symbol.dispose]() { + this.release(); + }, + }; + } + + request() { + if (!this.locked) { + return null; + } + + const { promise, resolve } = promiseWithResolvers(); + this.resolverQueue.push(resolve); + + return promise; + } +} + +export class CallSerializer2 { + private currentPromise: Promise | null = null; + private queuedCalls = 0; + + call(fn: () => T) { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + type ReturnType = T extends Promise ? T : T | Promise; + + if (this.currentPromise) { + this.queuedCalls++; + + return (this.currentPromise = this.currentPromise + .catch(() => {}) + .then(() => { + this.queuedCalls--; + return fn(); + }) + .finally(() => { + if (this.queuedCalls === 0) { + this.currentPromise = null; + } + })) as unknown as ReturnType; + } else { + const result = fn(); + + if (result instanceof Promise) { + this.currentPromise = result + .finally(() => { + if (this.queuedCalls === 0) { + this.currentPromise = null; + } + }); + } + + return result as unknown as ReturnType; + } + } + + done() { + if (this.currentPromise) { + return this.currentPromise + .catch(() => {}) + .then(() => {}); + } else { + return null; + } + } +} + +export const defer = (callback: () => void) => { + let executed = false; + + return { + execute() { + if (executed) { + return; + } + + executed = true; + callback(); + }, + [Symbol.dispose]() { + this.execute(); + }, + }; +}; + +export const promiseIterateAll = async function* (iterable: Iterable) { + for (const promise of iterable) { + yield await promise; + } +}; diff --git a/src/output.ts b/src/output.ts index 4046c1f..9362805 100644 --- a/src/output.ts +++ b/src/output.ts @@ -77,7 +77,7 @@ export type BaseTrackMetadata = { /** The track's disposition, i.e. information about its intended usage. */ disposition?: Partial; /** - * The maximum amount of encoded packets that will be added to this track. Setting this field provides the muxer + * The maximum number of encoded packets that will be added to this track. Setting this field provides the muxer * with an additional signal that it can use to preallocate space in the file. * * When this field is set, it is an error to provide more packets than whatever this field specifies. diff --git a/src/packet.ts b/src/packet.ts index 4f9b749..62e4efc 100644 --- a/src/packet.ts +++ b/src/packet.ts @@ -68,7 +68,7 @@ export class EncodedPacket { /** The duration of this packet in seconds. */ public readonly duration: number, /** - * The sequence number indicates the decode order of the packets. Packet A must be decoded before packet B if A + * The sequence number indicates the decode order of the packets. Packet A must be decoded before packet B if A * has a lower sequence number than B. If two packets have the same sequence number, they are the same packet. * Otherwise, sequence numbers are arbitrary and are not guaranteed to have any meaning besides their relative * ordering. Negative sequence numbers mean the sequence number is undefined. diff --git a/src/sample.ts b/src/sample.ts index ddd4527..b8b3be0 100644 --- a/src/sample.ts +++ b/src/sample.ts @@ -53,6 +53,9 @@ export type VideoSampleInit = { * @public */ export class VideoSample implements Disposable { + /** @internal */ + static _openSampleCount = 0; + /** @internal */ _data!: VideoFrame | OffscreenCanvas | Uint8Array | null; /** @internal */ @@ -107,6 +110,14 @@ export class VideoSample implements Disposable { return this.format && this.format.includes('A'); } + /** + * Whether this sample is closed, meaning its underlying data has been discarded. When a sample is closed, most + * operations will fail. + */ + get closed() { + return this._closed; + } + /** * Creates a new {@link VideoSample} from a * [`VideoFrame`](https://developer.mozilla.org/en-US/docs/Web/API/VideoFrame). This is essentially a near zero-cost @@ -265,6 +276,8 @@ export class VideoSample implements Disposable { } else { throw new TypeError('Invalid data type: Must be a BufferSource or CanvasImageSource.'); } + + VideoSample._openSampleCount++; } /** Clones this video sample. */ @@ -320,6 +333,7 @@ export class VideoSample implements Disposable { } this._closed = true; + VideoSample._openSampleCount--; } /** Returns the number of bytes required to hold this video sample's pixel data. */ @@ -888,6 +902,14 @@ export class AudioSample implements Disposable { return Math.trunc(SECOND_TO_MICROSECOND_FACTOR * this.duration); } + /** + * Whether this sample is closed, meaning its underlying data has been discarded. When a sample is closed, most + * operations will fail. + */ + get closed() { + return this._closed; + } + /** * Creates a new {@link AudioSample}, either from an existing * [`AudioData`](https://developer.mozilla.org/en-US/docs/Web/API/AudioData) or from raw bytes specified in diff --git a/test/browser/sample-cursor.test.ts b/test/browser/sample-cursor.test.ts new file mode 100644 index 0000000..60bbfdf --- /dev/null +++ b/test/browser/sample-cursor.test.ts @@ -0,0 +1,531 @@ +import { expect, test } from 'vitest'; +import { Input } from '../../src/input.js'; +import { 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 { promiseIterateAll } from '../../src/misc.js'; + +test('Sample cursor seeking', 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 cursor = await VideoSampleCursor.init(reader); + cursor.debugInfo.enabled = true; + + expect(cursor.current).toBe(null); + expect(cursor.debugInfo.pumpsStarted).toBe(0); + + // Seek to start + const seekToResult1 = cursor.seekToFirst(); + expect(seekToResult1).instanceOf(Promise); + const sample1 = (await seekToResult1)!; + expect(sample1).not.toBe(null); + expect(sample1).toBe(cursor.current); + expect(sample1.timestamp).toBeLessThanOrEqual(0); + expect(sample1.closed).toBe(false); + + expect(cursor.debugInfo.pumpsStarted).toBe(1); + + // Seek to a frame in the current GOP + const seekToResult2 = cursor.seekTo(0.5); + expect(seekToResult2).instanceOf(Promise); + const sample2 = (await seekToResult2)!; + expect(sample2).not.toBe(null); + expect(sample2).toBe(cursor.current); + expect(sample2.timestamp).toBeGreaterThan(0); + expect(sample2.timestamp).toBeLessThanOrEqual(0.5); + expect(sample2.closed).toBe(false); + expect(sample1.closed).toBe(true); + + // Wait a little bit so the next sample is most definitely decoded and waiting in the queue + await new Promise(resolve => setTimeout(resolve, 200)); + + const seekToResult3 = cursor.seekTo(0.55); + expect(seekToResult3).instanceOf(VideoSample); // No promise this time + const sample3 = (await seekToResult3)!; + expect(sample3).toBe(cursor.current); + expect(sample3.timestamp).toBeGreaterThan(0.5); + expect(sample3.timestamp).toBeLessThanOrEqual(0.55); + expect(sample3.closed).toBe(false); + expect(sample2.closed).toBe(true); + + // Let's get the same sample again + const seekToResult4 = cursor.seekTo(0.55); + expect(seekToResult4).instanceOf(VideoSample); + const sample4 = (await seekToResult4)!; + expect(sample3).toBe(sample4); + + expect(cursor.debugInfo.pumpsStarted).toBe(1); + + // Seek to a different GOP + const seekToResult5 = cursor.seekTo(2); + expect(seekToResult5).instanceOf(Promise); + const sample5 = (await seekToResult5)!; + expect(sample5).toBe(cursor.current); + expect(sample5.timestamp).toBeGreaterThan(0.55); + expect(sample5.timestamp).toBeLessThanOrEqual(2); + expect(sample5.closed).toBe(false); + expect(sample3.closed).toBe(true); + + expect(cursor.debugInfo.pumpsStarted).toBe(2); + + // Seek to a frame in the current GOP + const seekToResult6 = cursor.seekTo(2.5); + expect(seekToResult6).instanceOf(Promise); + const sample6 = (await seekToResult6)!; + expect(sample6).toBe(cursor.current); + expect(sample6.timestamp).toBeGreaterThan(2); + expect(sample6.timestamp).toBeLessThanOrEqual(2.5); + expect(sample6.closed).toBe(false); + expect(sample5.closed).toBe(true); + + // Seek to a frame in the current GOP, but backwards, requiring decoding to start over + const seekToResult7 = cursor.seekTo(2.4); + expect(seekToResult7).instanceOf(Promise); + const sample7 = (await seekToResult7)!; + expect(sample7).toBe(cursor.current); + expect(sample7.timestamp).toBeGreaterThan(2); + expect(sample7.timestamp).toBeLessThanOrEqual(2.4); + expect(sample7.closed).toBe(false); + expect(sample6.closed).toBe(true); + + expect(cursor.debugInfo.pumpsStarted).toBe(3); + + // Seek to a previous GOP + const seekToResult8 = cursor.seekTo(1); + expect(seekToResult8).instanceOf(Promise); + const sample8 = (await seekToResult8)!; + expect(sample8).toBe(cursor.current); + expect(sample8.timestamp).toBeLessThanOrEqual(1); + expect(sample8.closed).toBe(false); + expect(sample7.closed).toBe(true); + + expect(cursor.debugInfo.pumpsStarted).toBe(4); + + // Seek to past the end + const seekToResult9 = cursor.seekTo(Infinity); + expect(seekToResult9).instanceOf(Promise); + const sample9 = (await seekToResult9)!; + expect(sample9).toBe(cursor.current); + expect(sample9.timestamp).toBe(5); // The length of the video + expect(sample9.closed).toBe(false); + expect(sample8.closed).toBe(true); + + expect(cursor.debugInfo.pumpsStarted).toBe(5); + + // Seek to before the start + const seekToResult10 = cursor.seekTo(-Infinity); + expect(seekToResult10).toBe(null); + expect(sample9.closed).toBe(true); + + const seekToResult11 = cursor.seekToKey(2.5); + expect(seekToResult11).toBeInstanceOf(Promise); + const sample11 = (await seekToResult11)!; + expect(sample11).toBe(cursor.current); + expect(sample11.timestamp).toBe(2); + expect(sample11.closed).toBe(false); + + await cursor.close(); + + expect(sample11.closed).toBe(true); + expect(cursor.current).toBe(null); + expect(cursor.debugInfo.pumpsStarted).toBe(6); + + await cursor.close(); + + expect(VideoSample._openSampleCount).toBe(0); +}); + +test('Sample cursor advancing', 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 cursor = await VideoSampleCursor.init(reader); + cursor.debugInfo.enabled = true; + + expect(cursor.current).toBe(null); + expect(cursor.debugInfo.pumpsStarted).toBe(0); + + const firstSample = (await cursor.seekToFirst())!; + const secondSample = (await cursor.next())!; + const thirdSample = (await cursor.next())!; + + expect(secondSample.timestamp).toBeGreaterThan(firstSample.timestamp); + expect(thirdSample.timestamp).toBeGreaterThan(secondSample.timestamp); + expect(firstSample.closed).toBe(true); + expect(secondSample.closed).toBe(true); + expect(thirdSample.closed).toBe(false); + expect(cursor.current).toBe(thirdSample); + + await new Promise(resolve => setTimeout(resolve, 200)); + + const fourthSampleResult = cursor.next(); // It's available instantly + expect(fourthSampleResult).not.toBeInstanceOf(Promise); + + const fourthSample = fourthSampleResult as VideoSample; + expect(fourthSample.timestamp).toBeGreaterThan(thirdSample.timestamp); + + const lastSample = await cursor.seekTo(5); + expect(lastSample).not.toBe(null); + + const nextSample = await cursor.next(); + expect(nextSample).toBe(null); + const nextNextSample = await cursor.next(); + expect(nextNextSample).toBe(null); + + const middleSample = (await cursor.seekTo(3))!; + const sampleAfterMiddle = (await cursor.next())!; + expect(sampleAfterMiddle.timestamp).toBeGreaterThan(middleSample.timestamp); + + await cursor.seekTo(-Infinity); + const firstSampleAgain = (await cursor.next())!; + expect(firstSampleAgain.timestamp).toBe(firstSample.timestamp); + + let total = 0; + let lastTimestamp = -Infinity; + await cursor.iterate((sample) => { + total++; + + expect(sample.timestamp).toBeGreaterThan(lastTimestamp); + lastTimestamp = sample.timestamp; + }); + expect(total).toBe(121); + expect(cursor.current).toBe(null); + + total = 0; + await cursor.iterate(() => total++); + expect(total).toBe(0); // Since we're at the end + + await cursor.seekToFirst(); ; + total = 0; + await cursor.iterate((sample, stop) => { + if (sample.timestamp === 1) { + stop(); + return; + } + + total++; + }); + expect(total).toBe(24); + expect(cursor.current!.timestamp).toBe(1); + + total = 0; + for await (const sample of cursor) { + if (total === 0) { + expect(sample.timestamp).toBe(1); + } + + total++; + } + + expect(total).toBe(97); + + await cursor.close(); + + expect(cursor.debugInfo.pumpsStarted).toBe(5); + + expect(VideoSample._openSampleCount).toBe(0); +}); + +test('Sample cursor advancing, cold start', 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 cursor = await VideoSampleCursor.init(reader); + cursor.debugInfo.enabled = true; + + const firstSample = (await cursor.next())!; + expect(firstSample).not.toBe(null); + expect(firstSample.timestamp).toBe(0); + + const secondSample = (await cursor.next())!; + expect(secondSample.timestamp).toBeGreaterThan(firstSample.timestamp); + + await cursor.seekTo(-Infinity); + + void cursor.next(); + void cursor.seekTo(2); + + await cursor.close(); + + // 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); + for await (const sample of cursor2) { + expect(sample.timestamp).toBe(0); + break; + } + + await cursor2.iterate((sample, stop) => { + expect(sample.timestamp).toBe(0); + stop(); + }); + + await cursor2.close(); + + expect(VideoSample._openSampleCount).toBe(0); +}); + +test('Decoder pump error handling', 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 cursor = await VideoSampleCursor.init(reader); + cursor.debugInfo.enabled = true; + + cursor.debugInfo.throwInPump = true; + await expect(cursor.seekToFirst()).rejects.toThrow(); + + 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(); + + expect(VideoSample._openSampleCount).toBe(0); +}); + +test('Use after close', 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 cursor = await VideoSampleCursor.init(reader); + + await cursor.close(); + + 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(); +}); + +test('Command queuing', 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 cursor0 = await VideoSampleCursor.init(reader); + cursor0.debugInfo.enabled = true; + + const commands0 = [ + cursor0.seekToFirst(), + cursor0.seekToFirst(), + cursor0.close(), + ]; + expect(commands0.every(x => x instanceof Promise)).toBe(true); + + expect(cursor0.debugInfo.pumpsStarted).toBe(1); + + const cursor1 = await VideoSampleCursor.init(reader); + cursor1.debugInfo.enabled = true; + + const commands1 = [ + cursor1.seekToFirst(), + cursor1.close(), + ]; + expect(commands1.every(x => x instanceof Promise)).toBe(true); + + // eslint-disable-next-line @typescript-eslint/await-thenable + const results1 = await Promise.all(commands1); + expect(results1[0]!.timestamp).toBe(0); + expect(cursor1.debugInfo.decodedPackets.map(x => x.timestamp)).toEqual([0]); + + const cursor2 = await VideoSampleCursor.init(reader); + cursor2.debugInfo.enabled = true; + + const commands2 = [ + cursor2.seekTo(0), + cursor2.seekTo(1), + cursor2.seekTo(2), + cursor2.seekTo(3), + cursor2.seekTo(4), + cursor2.seekTo(5), + cursor2.close(), + ]; + expect(commands2.every(x => x instanceof Promise)).toBe(true); + + // eslint-disable-next-line @typescript-eslint/await-thenable + const results2 = await Promise.all(commands2); + expect(results2[0]!.timestamp).toBe(0); + expect(results2[1]!.timestamp).toBe(1); + expect(results2[2]!.timestamp).toBe(2); + expect(results2[3]!.timestamp).toBe(3); + expect(results2[4]!.timestamp).toBe(4); + expect(results2[5]!.timestamp).toBe(5); + expect(cursor2.debugInfo.decodedPackets.map(x => x.timestamp)).toEqual([ + 0, 1, 2, 3, 4, 5, + ]); + + const cursor3 = await VideoSampleCursor.init(reader); + cursor3.debugInfo.enabled = true; + + const commands3 = [ + cursor3.seekTo(0.5), + cursor3.seekTo(0.4), + cursor3.seekTo(0.3), + cursor3.seekTo(0.2), + cursor3.seekTo(0.1), + cursor3.seekTo(0), + cursor3.close(), + ]; + expect(commands3.every(x => x instanceof Promise)).toBe(true); + + // eslint-disable-next-line @typescript-eslint/await-thenable + const results3 = await Promise.all(commands3); + + expect(results3[0]!.timestamp).toBeLessThanOrEqual(0.5); + expect(results3[1]!.timestamp).toBeLessThanOrEqual(0.4); + expect(results3[2]!.timestamp).toBeLessThanOrEqual(0.3); + expect(results3[3]!.timestamp).toBeLessThanOrEqual(0.2); + expect(results3[4]!.timestamp).toBeLessThanOrEqual(0.1); + expect(results3[5]!.timestamp).toBe(0); + expect(cursor3.debugInfo.decodedPackets.every(x => x.timestamp <= 0.5)).toBe(true); + expect(cursor3.debugInfo.pumpsStarted).toBe(1); + + const cursor4 = await VideoSampleCursor.init(reader); + cursor4.debugInfo.enabled = true; + + const commands4 = [ + cursor4.seekToFirst(), + cursor4.next(), + cursor4.next(), + cursor4.close(), + ]; + + expect(commands4.every(x => x instanceof Promise)).toBe(true); + + // eslint-disable-next-line @typescript-eslint/await-thenable + const results4 = await Promise.all(commands4); + + expect(results4[0]!.timestamp).toBe(0); + expect(results4[1]!.timestamp).toBeGreaterThan(results4[0]!.timestamp); + 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); + cursor5.debugInfo.enabled = true; + + const commands5 = [ + cursor5.seekTo(0), + cursor5.next(), + cursor5.next(), + cursor5.seekTo(0.4), + cursor5.next(), + cursor5.next(), + cursor5.seekTo(0.8), + cursor5.next(), + cursor5.next(), + cursor5.seekTo(0), + cursor5.next(), + cursor5.close(), + ]; + + // eslint-disable-next-line @typescript-eslint/await-thenable + const results5 = await Promise.all(commands5); + + expect(results5[0]!.timestamp).toBe(0); + expect(results5[1]!.timestamp).toBeGreaterThan(results5[0]!.timestamp); + expect(results5[2]!.timestamp).toBeGreaterThan(results5[1]!.timestamp); + expect(results5[3]!.timestamp).toBeLessThanOrEqual(0.4); + expect(results5[3]!.timestamp).toBeGreaterThan(results5[2]!.timestamp); + expect(results5[4]!.timestamp).toBeGreaterThan(results5[3]!.timestamp); + expect(results5[5]!.timestamp).toBeGreaterThan(results5[4]!.timestamp); + expect(results5[6]!.timestamp).toBeLessThanOrEqual(0.8); + expect(results5[6]!.timestamp).toBeGreaterThan(results5[5]!.timestamp); + expect(results5[7]!.timestamp).toBeGreaterThan(results5[6]!.timestamp); + expect(results5[8]!.timestamp).toBeGreaterThan(results5[7]!.timestamp); + expect(results5[9]!.timestamp).toBe(0); + expect(results5[10]!.timestamp).toBeGreaterThan(results5[9]!.timestamp); + + expect(cursor5.debugInfo.pumpsStarted).toBe(1); + + const cursor6 = await VideoSampleCursor.init(reader); + cursor6.debugInfo.enabled = true; + + const commands6 = [ + cursor6.seekTo(0), + cursor6.seekTo(0.4), + cursor6.seekTo(0.8), + cursor6.seekTo(3.4), + cursor6.seekTo(3.8), + cursor6.seekTo(5), + cursor6.close(), + ]; + + // eslint-disable-next-line @typescript-eslint/await-thenable + const results6 = await Promise.all(commands6); + + expect(results6[0]!.timestamp).toBe(0); + expect(results6[1]!.timestamp).toBeLessThanOrEqual(0.4); + expect(results6[2]!.timestamp).toBeLessThanOrEqual(0.8); + expect(results6[3]!.timestamp).toBeLessThanOrEqual(3.4); + expect(results6[4]!.timestamp).toBeLessThanOrEqual(3.8); + expect(results6[5]!.timestamp).toBe(5); + + expect(cursor6.debugInfo.pumpsStarted).toBe(3); + + const cursor7 = await VideoSampleCursor.init(reader); + const commands7 = [ + cursor7.close(), + cursor7.close(), + ]; + expect(commands7[0]).toBe(commands7[1]); // Same Promise + + await Promise.all(commands7); + + const cursor8 = await VideoSampleCursor.init(reader, { autoClose: false }); + + const firstSample = await cursor8.seekToFirst(); + firstSample!.close(); + await new Promise(resolve => setTimeout(resolve, 200)); + + const commands8 = [ + cursor8.next(), + cursor8.next(), + cursor8.next(), + ]; + expect(commands8.every(x => !(x instanceof Promise))).toBe(true); + + for await (using sample of promiseIterateAll(commands8)) { + expect(sample!.closed).toBe(false); + } + + await cursor8.close(); + + 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 Sylvie video? diff --git a/test/node/call-serializer.test.ts b/test/node/call-serializer.test.ts new file mode 100644 index 0000000..7ff94e6 --- /dev/null +++ b/test/node/call-serializer.test.ts @@ -0,0 +1,58 @@ +import { expect, test } from 'vitest'; +import { CallSerializer2 } from '../../src/misc.js'; + +const executeDelayed = async (fn: () => T) => { + await new Promise(resolve => setTimeout(resolve, 10)); + return fn(); +}; + +test('Call serialization and return values', async () => { + const numbers: number[] = []; + const serializer = new CallSerializer2(); + + const first = serializer.call(() => numbers.push(1)); + const second = serializer.call(() => executeDelayed(() => numbers.push(2))); + const third = serializer.call(() => numbers.push(3)); + const fourth = serializer.call(() => executeDelayed(() => numbers.push(4))); + + await fourth; + + expect(numbers).toEqual([1, 2, 3, 4]); + expect(await first).toBe(1); + expect(await second).toBe(2); + expect(await third).toBe(3); + expect(await fourth).toBe(4); +}); + +test('Synchronous return value', async () => { + const serializer = new CallSerializer2(); + + const first = serializer.call(() => {}); + expect(first).not.toBeInstanceOf(Promise); + + const second = serializer.call(async () => {}); + const third = serializer.call(() => {}); + expect(second).toBeInstanceOf(Promise); + expect(third).toBeInstanceOf(Promise); + + await third; + + const fourth = serializer.call(() => {}); + expect(fourth).not.toBeInstanceOf(Promise); +}); + +test('Error handling', async () => { + const serializer = new CallSerializer2(); + + expect(() => serializer.call(() => { + throw new Error('yo'); + })).toThrow(); + + const second = serializer.call(() => executeDelayed(() => { + throw new Error('yo'); + })); + const third = serializer.call(() => executeDelayed(() => 1 + 2)); + + await expect(second).rejects.toThrow(); + expect(await third).toBe(3); +}); diff --git a/test/node/packet-cursor.test.ts b/test/node/packet-cursor.test.ts new file mode 100644 index 0000000..c95aff2 --- /dev/null +++ b/test/node/packet-cursor.test.ts @@ -0,0 +1,334 @@ +import { expect, test } from 'vitest'; +import { Input } from '../../src/input.js'; +import { BufferSource, FilePathSource } from '../../src/source.js'; +import path from 'node:path'; +import fs from 'node:fs'; +import { ALL_FORMATS } from '../../src/input-format.js'; +import { PacketCursor, PacketReader } from '../../src/cursors.js'; + +const __dirname = new URL('.', import.meta.url).pathname; + +test('Packet reader', async () => { + using input = new Input({ + source: new FilePathSource(path.join(__dirname, '../public/trim-buck-bunny.mov')), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + + const packet1 = (await reader.readFirst())!; + expect(packet1.timestamp).toBe(0); + + const packet3 = (await reader.readNext(packet1))!; + expect(packet3.sequenceNumber).toBeGreaterThan(packet1.sequenceNumber); + + const packet4 = (await reader.readNextKey(packet1))!; + expect(packet4.sequenceNumber).toBeGreaterThan(packet3.sequenceNumber); + expect(packet4.type).toBe('key'); + + const packet5 = (await reader.readNext(packet3))!; + expect(packet5.sequenceNumber).toBeGreaterThan(packet3.sequenceNumber); + expect(packet5.sequenceNumber).toBeLessThan(packet4.sequenceNumber); + + const packet6 = (await reader.readAt(2.4))!; + expect(packet6.timestamp).toBeGreaterThan(2); + expect(packet6.timestamp).toBeLessThanOrEqual(2.4); +}); + +test('Packet cursor seeking', async () => { + using input = new Input({ + source: new FilePathSource(path.join(__dirname, '../public/trim-buck-bunny.mov')), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + const cursor = new PacketCursor(reader); + + expect(cursor.current).toBe(null); + + const packet1 = (await cursor.seekToFirst())!; + expect(packet1).not.toBe(null); + expect(packet1).toBe(cursor.current); + expect(packet1.timestamp).toBe(0); + + const packet2 = (await cursor.seekTo(0.01))!; + expect(packet1.sequenceNumber).toBe(packet2.sequenceNumber); // Same packet + + const packet3 = (await cursor.seekTo(0.1))!; + expect(packet3).toBe(cursor.current); + expect(packet3.timestamp).toBeGreaterThan(0); + expect(packet3.sequenceNumber).toBeGreaterThan(packet1.sequenceNumber); + + const packet4 = (await cursor.seekToKey(0.1))!; + expect(packet4).toBe(cursor.current); + expect(packet4.sequenceNumber).toBe(packet1.sequenceNumber); + + const packet5 = (await cursor.seekTo(Infinity))!; + expect(packet5).toBe(cursor.current); + expect(packet5.timestamp).toBe(5); + + const packet6 = (await cursor.seekTo(-Infinity))!; + expect(packet6).toBe(cursor.current); + expect(packet6).toBe(null); +}); + +test('Packet cursor iteration', async () => { + using input = new Input({ + source: new FilePathSource(path.join(__dirname, '../public/trim-buck-bunny.mov')), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + const cursor = new PacketCursor(reader); + + const packet0 = (await cursor.seekToFirst())!; + expect(cursor.current!.timestamp).toBe(0); + + const packet1 = (await cursor.next())!; + expect(packet1.sequenceNumber).toBeGreaterThan(packet0.sequenceNumber); + expect(packet1).toBe(cursor.current); + + const packet2 = (await cursor.next())!; + expect(packet2.sequenceNumber).toBeGreaterThan(packet1.sequenceNumber); + expect(packet2).toBe(cursor.current); + + const packet3 = (await cursor.nextKey())!; + expect(packet3.sequenceNumber).toBeGreaterThan(packet2.sequenceNumber); + expect(packet3.type).toBe('key'); + expect(packet3).toBe(cursor.current); + + await cursor.seekTo(Infinity); + expect(cursor.current).not.toBe(null); + + const packet4 = await cursor.next(); + expect(packet4).toBe(null); + expect(packet4).toBe(cursor.current); + + const packet5 = await cursor.next(); + expect(packet5).toBe(null); + + await cursor.seekTo(-Infinity); + expect(cursor.current).toBe(null); + + const packet6 = (await cursor.next())!; + expect(packet6.sequenceNumber).toBe(packet0.sequenceNumber); + expect(packet6).toBe(cursor.current); + + await cursor.seekTo(-Infinity); + expect(cursor.current).toBe(null); + + const packet7 = (await cursor.next())!; + expect(packet7.sequenceNumber).toBe(packet0.sequenceNumber); + expect(packet7).toBe(cursor.current); + + const packet8 = (await cursor.next())!; + expect(packet8.sequenceNumber).toBeGreaterThan(packet7.sequenceNumber); + expect(packet8).toBe(cursor.current); + + await cursor.seekToFirst(); + + let total = 0; + let lastSeqNum = -Infinity; + for await (const packet of cursor) { + if (total === 0) { + expect(packet.sequenceNumber).toBe(packet0.sequenceNumber); + } + + expect(packet.sequenceNumber).toBeGreaterThan(lastSeqNum); + + lastSeqNum = packet.sequenceNumber; + total++; + } + + expect(total).toBe(121); + + for await (const _ of cursor) { + throw new Error('Unreachable'); + } + + total = 0; + await cursor.seekTo(1); + for await (const _ of cursor) { + total++; + } + + expect(total).toBe(97); + + total = 0; + await cursor.seekToFirst(); + await cursor.iterate(() => total++); + + expect(total).toBe(121); + expect(cursor.current).toBe(null); + + total = 0; + await cursor.seekToFirst(); + await cursor.iterate((packet, stop) => { + if (packet.timestamp === 1) { + stop(); + return; + } + + total++; + }); + + expect(total).toBe(24); + expect(cursor.current!.timestamp).toBe(1); + + await cursor.seekTo(-Infinity); + total = 0; + for await (const _ of cursor) total++; + + expect(total).toBe(121); + + await cursor.seekTo(-Infinity); + total = 0; + await cursor.iterate(() => total++); + + expect(total).toBe(121); + + const cursor2 = new PacketCursor(reader); + const packet9 = (await cursor2.next())!; // Without any prior seeks + expect(packet9.sequenceNumber).toBe(packet0.sequenceNumber); +}); + +test('Synchronous packet reading', async () => { + using input = new Input({ + source: new BufferSource(fs.readFileSync(path.join(__dirname, '../public/trim-buck-bunny.mov'))), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + const cursor = new PacketCursor(reader); + + expect(reader.readFirst()).not.toBeInstanceOf(Promise); + + expect(cursor.seekToFirst()).not.toBeInstanceOf(Promise); + expect(cursor.seekTo(0.1)).not.toBeInstanceOf(Promise); + expect(cursor.seekToKey(0.1)).not.toBeInstanceOf(Promise); + expect(cursor.seekTo(2)).not.toBeInstanceOf(Promise); + expect(cursor.seekTo(Infinity)).not.toBeInstanceOf(Promise); + expect(cursor.seekTo(-Infinity)).not.toBeInstanceOf(Promise); + + void cursor.seekToFirst(); + + expect(cursor.next()).not.toBeInstanceOf(Promise); +}); + +test('Command queuing', async () => { + using input = new Input({ + source: new FilePathSource(path.join(__dirname, '../public/trim-buck-bunny.mov'), { + maxCacheSize: 0, // So all commands return promises + }), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + const cursor = new PacketCursor(reader); + + expect(cursor.waitUntilIdle()).toBe(null); + + const commands = [ + cursor.seekToFirst(), + cursor.next(), + cursor.next(), + cursor.seekTo(2.4), + cursor.next(), + cursor.nextKey(), + cursor.waitUntilIdle()!.then(() => cursor.current), + cursor.nextKey(), + cursor.nextKey(), + cursor.next(), + cursor.seekTo(Infinity), + cursor.seekTo(-Infinity), + cursor.seekToKey(2.4), + ]; + + expect(commands.every(x => x instanceof Promise)).toBe(true); + + // eslint-disable-next-line @typescript-eslint/await-thenable + const resolved = await Promise.all(commands); + + expect(resolved[0]!.timestamp).toBe(0); + + expect(resolved[1]!.sequenceNumber).toBeGreaterThan(resolved[0]!.sequenceNumber); + + expect(resolved[2]!.sequenceNumber).toBeGreaterThan(resolved[1]!.sequenceNumber); + + expect(resolved[3]!.timestamp).toBeGreaterThan(2); + expect(resolved[3]!.timestamp).toBeLessThanOrEqual(2.4); + + expect(resolved[4]!.sequenceNumber).toBeGreaterThan(resolved[3]!.sequenceNumber); + + expect(resolved[5]!.timestamp).toBe(3); + + expect(resolved[6]!.sequenceNumber).toBe(resolved[5]!.sequenceNumber); + + expect(resolved[7]!.timestamp).toBe(4); + + expect(resolved[8]!.timestamp).toBe(5); + + expect(resolved[9]).toBe(null); + + expect(resolved[10]!.sequenceNumber).toBe(resolved[8]!.sequenceNumber); + + expect(resolved[11]).toBe(null); + + expect(resolved[12]!.timestamp).toBe(2); + expect(resolved[12]!.type).toBe('key'); + + void cursor.seekTo(1); + await cursor.iterate((packet, stop) => { + expect(packet.timestamp).toBe(1); + stop(); + }); + + void cursor.seekTo(3); + for await (const packet of cursor) { + expect(packet.timestamp).toBe(3); + break; + } +}); + +test('verifyKeyPackets with faultily-labeled key frames', async () => { + using input = new Input({ + source: new FilePathSource(path.join(__dirname, '../public/fake-cod.mp4')), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + + const firstPacket = (await reader.readFirst())!; + expect(firstPacket.type).toBe('key'); + + const fakeKeyPacket = (await reader.readNextKey(firstPacket))!; + expect(fakeKeyPacket).not.toBe(null); + expect(fakeKeyPacket.type).toBe('key'); // Metadata says it's a key frame + expect(fakeKeyPacket.sequenceNumber).toBeGreaterThan(firstPacket.sequenceNumber); + + const verifiedPacket = (await reader.readAt(fakeKeyPacket.timestamp, { verifyKeyPackets: true }))!; + expect(verifiedPacket.sequenceNumber).toBe(fakeKeyPacket.sequenceNumber); + expect(verifiedPacket.type).toBe('delta'); // After verification, it's actually a delta frame + + const unverifiedKeyAt = (await reader.readKeyAt(fakeKeyPacket.timestamp))!; + expect(unverifiedKeyAt.sequenceNumber).toBe(fakeKeyPacket.sequenceNumber); + expect(unverifiedKeyAt.type).toBe('key'); + + const verifiedKeyAt = (await reader.readKeyAt(fakeKeyPacket.timestamp, { verifyKeyPackets: true }))!; + expect(verifiedKeyAt.sequenceNumber).toBe(firstPacket.sequenceNumber); + expect(verifiedKeyAt.type).toBe('key'); + + const unverifiedNextKey = (await reader.readNextKey(firstPacket))!; + expect(unverifiedNextKey).not.toBe(null); + expect(unverifiedNextKey.type).toBe('key'); + expect(unverifiedNextKey.sequenceNumber).toBe(fakeKeyPacket.sequenceNumber); + + const verifiedNextKey = await reader.readNextKey(firstPacket, { verifyKeyPackets: true }); + expect(verifiedNextKey).toBe(null); +}); diff --git a/test/public/fake-cod.mp4 b/test/public/fake-cod.mp4 new file mode 100644 index 0000000..901c608 Binary files /dev/null and b/test/public/fake-cod.mp4 differ diff --git a/test/public/trim-buck-bunny.mov b/test/public/trim-buck-bunny.mov new file mode 100644 index 0000000..4780a5e Binary files /dev/null and b/test/public/trim-buck-bunny.mov differ