From f045e5988c83b587732b20fa95c0d57ef00f8b06 Mon Sep 17 00:00:00 2001 From: Vanilagy <1696106+Vanilagy@users.noreply.github.com> Date: Mon, 15 Dec 2025 17:07:22 +0100 Subject: [PATCH] Implement sample cursor reset() --- src/cursors.ts | 88 ++++++++++++++++-------- test/browser/sample-cursor.test.ts | 105 +++++++++++++++++++++++++++-- 2 files changed, 160 insertions(+), 33 deletions(-) diff --git a/src/cursors.ts b/src/cursors.ts index eb1357c..cb63151 100644 --- a/src/cursors.ts +++ b/src/cursors.ts @@ -343,7 +343,7 @@ type PendingRequest = { successor: PendingRequest | null; }; -type SampleTransformer = (sample: Sample) => MaybePromise; +type SampleTransformer = (sample: Sample) => TransformedSample; type SampleCursorOptions = { autoClose?: boolean; @@ -379,7 +379,6 @@ export abstract class SampleCursor< pumpTarget: EncodedPacket | null = null; pumpMutex = new AsyncMutex4(); _closed = false; - otherMutex = new AsyncMutex4(); error: unknown = null; errorSet = false; @@ -423,15 +422,17 @@ export abstract class SampleCursor< .finally(() => lock.release()); } - async onDecoderSample(sample: Sample) { + onDecoderSample(sample: Sample): void { try { if (this.debugInfo.enabled && this.debugInfo.throwDecoderError) { sample.close(); - return this.onDecoderError(new Error('Fake decoder error!')); - } - using lock = this.otherMutex.lock(); - if (lock.pending) await lock.ready; + if (!this._closed) { + return this.onDecoderError(new Error('Fake decoder error!')); + } else { + return; + } + } while (this.decodedTimestamps.length > 0 && this.decodedTimestamps[0]! <= sample.timestamp) { this.decodedTimestamps.shift(); @@ -453,12 +454,7 @@ export abstract class SampleCursor< break; } - if (!transformed) { - let result = this.transform(sample); - if (result instanceof Promise) result = await result; - - transformed = result; - } + transformed ??= this.transform(sample); request.resolve(transformed); this.pendingRequests.splice(i--, 1); @@ -479,15 +475,12 @@ export abstract class SampleCursor< this.queueDequeue.resolve(); this.queueDequeue = promiseWithResolvers(); } catch (error) { - await this.closeWithError(error); + void this.closeWithError(error); } } - async onDecoderError(error: unknown) { - using lock = this.otherMutex.lock(); - if (lock.pending) await lock.ready; - - await this.closeWithError(error); + onDecoderError(error: unknown): void { + void this.closeWithError(error); } onDecoderDequeue() { @@ -569,9 +562,7 @@ export abstract class SampleCursor< this.queueDequeue = promiseWithResolvers(); if (targetPacket.timestamp <= nextSample.timestamp) { - let transformed = this.transform(nextSample); - if (transformed instanceof Promise) transformed = await transformed; - + const transformed = this.transform(nextSample); this.setCurrent(nextSample, transformed); return res.set(transformed); } else { @@ -720,9 +711,7 @@ export abstract class SampleCursor< this.queueDequeue.resolve(); this.queueDequeue = promiseWithResolvers(); - let transformed = this.transform(nextSample); - if (transformed instanceof Promise) transformed = await transformed; - + const transformed = this.transform(nextSample); this.setCurrent(nextSample, transformed); return res.set(transformed); } @@ -853,10 +842,13 @@ export abstract class SampleCursor< closePromise: Promise | null = null; close() { - return this.closePromise ??= this.closeInternal(); + return this.closePromise ??= this._closed + ? Promise.resolve() + : this.closeInternal(); } async closeInternal(doLock = true) { + // Almost correct, but not quite this.predictedRequests++; this.packetReader.track.input._openSampleCursors.delete(this); @@ -875,6 +867,7 @@ export abstract class SampleCursor< this.setCurrent(null, null); this.decoder?.close(); + this.decoder = null; } [Symbol.asyncDispose]() { @@ -998,10 +991,49 @@ export abstract class SampleCursor< return this.closeInternal(doLock); } - closeWithErrorAndThrow(error: unknown): never { - void this.closeWithError(error); + closeWithErrorAndThrow(error: unknown, doLock?: boolean): never { + void this.closeWithError(error, doLock); throw error; } + + async reset() { + this.predictedRequests++; + + using lock = this.pumpMutex.lock(); + if (lock.pending) await lock.ready; + + if (!this._closed) { + await this.closeInternal(false); + } + + assert(!this.pumpRunning); + assert(!this.pumpStopQueued); + assert(!this.currentRaw); + assert(!this.current); + assert(this.sampleQueue.length === 0); + assert(this.pendingRequests.length === 0); + assert(this.lastPendingRequest === null); + assert(this.decodedTimestamps.length === 0); + assert(this.maxDecodedSequenceNumber === -1); + assert(this.pumpTarget === null); + assert(!this.decoder || this.decoder.closed); + + this._closed = false; + this.closePromise = null; + this.error = null; + this.errorSet = false; + this.nextIsFirst = true; + this.predictedRequests = 0; + + this.packetReader.track.input._openSampleCursors.add(this); + + try { + const newDecoder = await this.initDecoder(); + this.decoder = newDecoder; + } catch (error) { + this.closeWithErrorAndThrow(error); + } + } } export class VideoSampleCursor extends SampleCursor { diff --git a/test/browser/sample-cursor.test.ts b/test/browser/sample-cursor.test.ts index bc46a87..f39c167 100644 --- a/test/browser/sample-cursor.test.ts +++ b/test/browser/sample-cursor.test.ts @@ -286,7 +286,72 @@ test('Sample cursor advancing, cold start', async () => { expect(VideoSample._openSampleCount).toBe(0); }); -test('Decoder setup error', async () => { +test('Sample cursor reset', async () => { + using input = new Input({ + source: new UrlSource('/trim-buck-bunny.mov'), + formats: ALL_FORMATS, + }); + + const videoTrack = (await input.getPrimaryVideoTrack())!; + const reader = new PacketReader(videoTrack); + await using cursor = new VideoSampleCursor(reader); + + expect(cursor.closed).toBe(false); + + await cursor.close(); + expect(cursor.closed).toBe(true); + + expect(() => cursor.seekToFirst()).toThrow('cursor has been closed'); + + await cursor.reset(); + expect(cursor.closed).toBe(false); + + const firstSample = await cursor.seekToFirst(); + expect(firstSample).not.toBe(null); + + await cursor.close(); + await cursor.reset(); + + const nextSample = await cursor.next(); + expect(nextSample).not.toBe(null); + expect(nextSample!.timestamp).toBe(firstSample!.timestamp); + + // Absolutely deranged usage, but it's gotta work! + const commands = [ + cursor.next(), + cursor.close(), + cursor.reset(), + cursor.next(), + cursor.next(), + cursor.reset(), + cursor.next(), + ]; + + // eslint-disable-next-line @typescript-eslint/await-thenable + const results = await Promise.all(commands); + + expect(cursor.closed).toBe(false); + expect(results[0]!.timestamp).toBeGreaterThan(firstSample!.timestamp); + expect(results[3]!.timestamp).toBe(firstSample!.timestamp); + expect(results[4]!.timestamp).toBe(results[0]!.timestamp); + expect(results[6]!.timestamp).toBe(firstSample!.timestamp); + + await using cursor2 = new VideoSampleCursor(reader); + cursor2.debugInfo.enabled = true; + + // Test if queueing a reset makes the decoder decode minimally many packets + const commands2 = [ + cursor2.seekToFirst(), + cursor2.reset(), + ]; + + // eslint-disable-next-line @typescript-eslint/await-thenable + await Promise.all(commands2); + + expect(cursor2.debugInfo.decodedPackets).toHaveLength(1); +}); + +test('Decoder setup error & reset', async () => { using input = new Input({ source: new UrlSource('/trim-buck-bunny.mov'), formats: ALL_FORMATS, @@ -297,18 +362,32 @@ test('Decoder setup error', async () => { await using cursor1 = new VideoSampleCursor(reader); cursor1.debugInfo.enabled = true; cursor1.debugInfo.throwInDecoderInit = true; + expect(cursor1.closed).toBe(false); await expect(cursor1.seekToFirst()).rejects.toThrow('Fake decoder init error'); expect(cursor1.closed).toBe(true); + expect(() => cursor1.seekToFirst()).toThrow('Fake decoder init error'); // Bricked + + await expect(cursor1.reset()).rejects.toThrow('Fake decoder init error'); + expect(cursor1.closed).toBe(true); + + cursor1.debugInfo.throwInDecoderInit = false; + await cursor1.reset(); + expect(cursor1.closed).toBe(false); + + const firstSample = (await cursor1.seekToFirst())!; + expect(firstSample).not.toBe(null); + + // Let's test directly closing after opening const cursor2 = new VideoSampleCursor(reader); cursor2.debugInfo.enabled = true; cursor2.debugInfo.throwInDecoderInit = true; await cursor2.close(); }); -test('Decoder pump error handling', async () => { +test('Decoder pump error handling & reset', async () => { using input = new Input({ source: new UrlSource('/trim-buck-bunny.mov'), formats: ALL_FORMATS, @@ -320,18 +399,23 @@ test('Decoder pump error handling', async () => { cursor.debugInfo.enabled = true; cursor.debugInfo.throwInPump = true; - await expect(async () => cursor.seekToFirst()).rejects.toThrow(); + await expect(async () => cursor.seekToFirst()).rejects.toThrow('Fake pump error'); expect(cursor.closed).toBe(true); expect(cursor.current).toBe(null); cursor.debugInfo.throwInPump = false; - await expect(async () => cursor.seekToFirst()).rejects.toThrow(); // It's bricked + await expect(async () => cursor.seekToFirst()).rejects.toThrow('Fake pump error'); // It's bricked expect(VideoSample._openSampleCount).toBe(0); + + await cursor.reset(); + + const firstSample = await cursor.seekToFirst(); + expect(firstSample!.timestamp).toBe(0); }); -test('Decoder errors', async () => { +test('Decoder errors & reset', async () => { using input = new Input({ source: new UrlSource('/trim-buck-bunny.mov'), formats: ALL_FORMATS, @@ -347,6 +431,14 @@ test('Decoder errors', async () => { await expect(cursor1.seekToFirst()).rejects.toThrow('Fake decoder error'); expect(cursor1.closed).toBe(true); + cursor1.debugInfo.throwDecoderError = false; + expect(() => cursor1.seekToFirst()).toThrow('Fake decoder error'); // Bricked + + await cursor1.reset(); + + const firstSample = await cursor1.seekToFirst(); + expect(firstSample!.timestamp).toBe(0); + const cursor2 = new VideoSampleCursor(reader); cursor2.debugInfo.enabled = true; @@ -392,6 +484,9 @@ test('Use after close', async () => { await expect(commands2[2]).resolves.toBeUndefined(); await expect(commands2[3]).rejects.toThrow('cursor has been closed'); + expect(cursor2.pumpRunning).toBe(false); + expect(cursor2.decoder).toBe(null); + const cursor3 = new VideoSampleCursor(reader); const commands3 = [ cursor3.seekToFirst(),