From 90ec40166fcb7a8ba78bc9ad9e70184b8e0df05c Mon Sep 17 00:00:00 2001 From: Vanilagy <1696106+Vanilagy@users.noreply.github.com> Date: Wed, 29 Jul 2026 11:42:14 +0200 Subject: [PATCH] Fix decoders not being closed in certain error scenarios, fixed unsettled promises in reading logic (fixes #446) --- src/media-sink.ts | 16 +++++++++++----- src/source.ts | 28 +++++++++++++++++++++++----- 2 files changed, 34 insertions(+), 10 deletions(-) diff --git a/src/media-sink.ts b/src/media-sink.ts index 9a57bb6..aa80639 100644 --- a/src/media-sink.ts +++ b/src/media-sink.ts @@ -471,6 +471,7 @@ export abstract class BaseMediaSampleSink< let decoderIsFlushed = false; let ended = false; let terminated = false; + let decoder: DecoderWrapper | null = null; // This stores errors that are "out of band" in the sense that they didn't occur in the normal flow of this // method but instead in a different context. This error should not go unnoticed and must be bubbled up to @@ -486,7 +487,7 @@ export abstract class BaseMediaSampleSink< // The following is the "pump" process that keeps pumping packets into the decoder (async () => { - const decoder = await this._createDecoder((sample) => { + decoder = await this._createDecoder((sample) => { onQueueDequeue(); if (sample.timestamp >= endTimestamp) { ended = true; @@ -569,20 +570,21 @@ export abstract class BaseMediaSampleSink< if (!terminated && !this._track.input._disposed) { await decoder.flush(); } - decoder.close(); if (!firstSampleQueued && lastSample) { sampleQueue.push(lastSample); } decoderIsFlushed = true; - onQueueNotEmpty(); // To unstuck the generator + onQueueNotEmpty(); // To unstuck (unstick?) the generator })().catch((error) => { if (!hasOutOfBandError) { outOfBandError = error; hasOutOfBandError = true; onQueueNotEmpty(); } + }).finally(() => { + decoder?.close(); }); const track = this._track; @@ -647,6 +649,7 @@ export abstract class BaseMediaSampleSink< let { promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers(); let decoderIsFlushed = false; let terminated = false; + let decoder: DecoderWrapper | null = null; // This stores errors that are "out of band" in the sense that they didn't occur in the normal flow of this // method but instead in a different context. This error should not go unnoticed and must be bubbled up to @@ -668,7 +671,7 @@ export abstract class BaseMediaSampleSink< // The following is the "pump" process that keeps pumping packets into the decoder (async () => { - const decoder = await this._createDecoder((sample) => { + decoder = await this._createDecoder((sample) => { onQueueDequeue(); if (terminated) { @@ -711,6 +714,7 @@ export abstract class BaseMediaSampleSink< const decodePackets = async () => { assert(lastKeyPacket); + assert(decoder); // Start at the current key packet let currentPacket = lastKeyPacket; @@ -738,6 +742,7 @@ export abstract class BaseMediaSampleSink< }; const flushDecoder = async () => { + assert(decoder); await decoder.flush(); // We don't expect this list to have any elements in it anymore, but in case it does, let's emit @@ -796,7 +801,6 @@ export abstract class BaseMediaSampleSink< await flushDecoder(); } - decoder.close(); decoderIsFlushed = true; onQueueNotEmpty(); // To unstuck the generator @@ -806,6 +810,8 @@ export abstract class BaseMediaSampleSink< hasOutOfBandError = true; onQueueNotEmpty(); } + }).finally(() => { + decoder?.close(); }); const track = this._track; diff --git a/src/source.ts b/src/source.ts index 4192d9b..cfa7cdb 100644 --- a/src/source.ts +++ b/src/source.ts @@ -1665,6 +1665,10 @@ export class ReadableStreamSource extends Source { /** @internal */ _dispose() { + for (const pendingSlice of this._pendingSlices) { + pendingSlice.reject(new InputDisposedError()); + } + this._pendingSlices.length = 0; this._cache.length = 0; void this._reader?.cancel(); @@ -2328,11 +2332,13 @@ class ReadOrchestrator { signalWorkerStoppedRunning(worker: ReadWorker) { worker.running = false; - // When a worker stops running, that means it has hit its targetPos. It might still have pendingSlices assigned, - // but this is because those pending slices cover data that other workers are assigned to fill. Since targetPos - // has been reached, we can confidently say that this worker has completed its share of work on the pending - // slices and must no longer care about them. - worker.pendingSlices.length = 0; + if (!worker.aborted) { + // When a worker stops running, that means it has hit its targetPos. It might still have pendingSlices + // assigned, but this is because those pending slices cover data that other workers are assigned to fill. + // Since targetPos has been reached, we can confidently say that this worker has completed its share of work + // on the pending slices and must no longer care about them. + worker.pendingSlices.length = 0; + } } /** Called when a worker reaches the end of the underlying data and must be cleaned up. */ @@ -2447,11 +2453,23 @@ class ReadOrchestrator { dispose() { for (const worker of this.workers) { + for (const slice of worker.pendingSlices) { + slice.reject(new InputDisposedError()); + } + + worker.pendingSlices.length = 0; worker.aborted = true; } + for (const queuedRead of this.queuedReads) { + for (const slice of queuedRead.pendingSlices) { + slice.reject(new InputDisposedError()); + } + } + this.workers.length = 0; this.cache.length = 0; + this.queuedReads.length = 0; this.disposed = true; } }