Fix decoders not being closed in certain error scenarios, fixed unsettled promises in reading logic (fixes #446)

This commit is contained in:
Vanilagy
2026-07-29 11:42:14 +02:00
parent 4b5ef55cf1
commit 90ec40166f
2 changed files with 34 additions and 10 deletions
+11 -5
View File
@@ -471,6 +471,7 @@ export abstract class BaseMediaSampleSink<
let decoderIsFlushed = false; let decoderIsFlushed = false;
let ended = false; let ended = false;
let terminated = false; let terminated = false;
let decoder: DecoderWrapper<MediaSample> | null = null;
// This stores errors that are "out of band" in the sense that they didn't occur in the normal flow of this // 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 // 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 // The following is the "pump" process that keeps pumping packets into the decoder
(async () => { (async () => {
const decoder = await this._createDecoder((sample) => { decoder = await this._createDecoder((sample) => {
onQueueDequeue(); onQueueDequeue();
if (sample.timestamp >= endTimestamp) { if (sample.timestamp >= endTimestamp) {
ended = true; ended = true;
@@ -569,20 +570,21 @@ export abstract class BaseMediaSampleSink<
if (!terminated && !this._track.input._disposed) { if (!terminated && !this._track.input._disposed) {
await decoder.flush(); await decoder.flush();
} }
decoder.close();
if (!firstSampleQueued && lastSample) { if (!firstSampleQueued && lastSample) {
sampleQueue.push(lastSample); sampleQueue.push(lastSample);
} }
decoderIsFlushed = true; decoderIsFlushed = true;
onQueueNotEmpty(); // To unstuck the generator onQueueNotEmpty(); // To unstuck (unstick?) the generator
})().catch((error) => { })().catch((error) => {
if (!hasOutOfBandError) { if (!hasOutOfBandError) {
outOfBandError = error; outOfBandError = error;
hasOutOfBandError = true; hasOutOfBandError = true;
onQueueNotEmpty(); onQueueNotEmpty();
} }
}).finally(() => {
decoder?.close();
}); });
const track = this._track; const track = this._track;
@@ -647,6 +649,7 @@ export abstract class BaseMediaSampleSink<
let { promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers(); let { promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers();
let decoderIsFlushed = false; let decoderIsFlushed = false;
let terminated = false; let terminated = false;
let decoder: DecoderWrapper<MediaSample> | null = null;
// This stores errors that are "out of band" in the sense that they didn't occur in the normal flow of this // 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 // 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 // The following is the "pump" process that keeps pumping packets into the decoder
(async () => { (async () => {
const decoder = await this._createDecoder((sample) => { decoder = await this._createDecoder((sample) => {
onQueueDequeue(); onQueueDequeue();
if (terminated) { if (terminated) {
@@ -711,6 +714,7 @@ export abstract class BaseMediaSampleSink<
const decodePackets = async () => { const decodePackets = async () => {
assert(lastKeyPacket); assert(lastKeyPacket);
assert(decoder);
// Start at the current key packet // Start at the current key packet
let currentPacket = lastKeyPacket; let currentPacket = lastKeyPacket;
@@ -738,6 +742,7 @@ export abstract class BaseMediaSampleSink<
}; };
const flushDecoder = async () => { const flushDecoder = async () => {
assert(decoder);
await decoder.flush(); await decoder.flush();
// We don't expect this list to have any elements in it anymore, but in case it does, let's emit // 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(); await flushDecoder();
} }
decoder.close();
decoderIsFlushed = true; decoderIsFlushed = true;
onQueueNotEmpty(); // To unstuck the generator onQueueNotEmpty(); // To unstuck the generator
@@ -806,6 +810,8 @@ export abstract class BaseMediaSampleSink<
hasOutOfBandError = true; hasOutOfBandError = true;
onQueueNotEmpty(); onQueueNotEmpty();
} }
}).finally(() => {
decoder?.close();
}); });
const track = this._track; const track = this._track;
+23 -5
View File
@@ -1665,6 +1665,10 @@ export class ReadableStreamSource extends Source {
/** @internal */ /** @internal */
_dispose() { _dispose() {
for (const pendingSlice of this._pendingSlices) {
pendingSlice.reject(new InputDisposedError());
}
this._pendingSlices.length = 0; this._pendingSlices.length = 0;
this._cache.length = 0; this._cache.length = 0;
void this._reader?.cancel(); void this._reader?.cancel();
@@ -2328,11 +2332,13 @@ class ReadOrchestrator {
signalWorkerStoppedRunning(worker: ReadWorker) { signalWorkerStoppedRunning(worker: ReadWorker) {
worker.running = false; worker.running = false;
// When a worker stops running, that means it has hit its targetPos. It might still have pendingSlices assigned, if (!worker.aborted) {
// but this is because those pending slices cover data that other workers are assigned to fill. Since targetPos // When a worker stops running, that means it has hit its targetPos. It might still have pendingSlices
// has been reached, we can confidently say that this worker has completed its share of work on the pending // assigned, but this is because those pending slices cover data that other workers are assigned to fill.
// slices and must no longer care about them. // Since targetPos has been reached, we can confidently say that this worker has completed its share of work
worker.pendingSlices.length = 0; // 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. */ /** Called when a worker reaches the end of the underlying data and must be cleaned up. */
@@ -2447,11 +2453,23 @@ class ReadOrchestrator {
dispose() { dispose() {
for (const worker of this.workers) { for (const worker of this.workers) {
for (const slice of worker.pendingSlices) {
slice.reject(new InputDisposedError());
}
worker.pendingSlices.length = 0;
worker.aborted = true; worker.aborted = true;
} }
for (const queuedRead of this.queuedReads) {
for (const slice of queuedRead.pendingSlices) {
slice.reject(new InputDisposedError());
}
}
this.workers.length = 0; this.workers.length = 0;
this.cache.length = 0; this.cache.length = 0;
this.queuedReads.length = 0;
this.disposed = true; this.disposed = true;
} }
} }