By the creators of AsyncMutex3

This commit is contained in:
Vanilagy
2025-12-14 23:32:55 +01:00
parent 0f015d8e98
commit fbcb4dcb95
2 changed files with 41 additions and 45 deletions
+11 -16
View File
@@ -4,7 +4,7 @@
import { PCM_AUDIO_CODECS } from './codec'; import { PCM_AUDIO_CODECS } from './codec';
import { InputAudioTrack, InputTrack, InputVideoTrack } from './input-track'; import { InputAudioTrack, InputTrack, InputVideoTrack } from './input-track';
import { AudioDecoderWrapper, DecoderWrapper, PacketRetrievalOptions, PcmAudioDecoderWrapper, validatePacketRetrievalOptions, validateTimestamp, VideoDecoderWrapper } from './media-sink'; import { AudioDecoderWrapper, DecoderWrapper, PacketRetrievalOptions, PcmAudioDecoderWrapper, validatePacketRetrievalOptions, validateTimestamp, VideoDecoderWrapper } from './media-sink';
import { assert, AsyncMutex3, AsyncMutexLock, CallSerializer2, defer, insertSorted, isFirefox, last, MaybePromise, polyfillSymbolDispose, promiseWithResolvers, ResultValue, Rotation, Yo } from './misc'; import { assert, AsyncMutex4, AsyncMutexLock, CallSerializer2, defer, insertSorted, isFirefox, last, MaybePromise, polyfillSymbolDispose, promiseWithResolvers, ResultValue, Rotation, Yo } from './misc';
import { EncodedPacket } from './packet'; import { EncodedPacket } from './packet';
import { AudioSample, clampCropRectangle, CropRectangle, validateCropRectangle, VideoSample } from './sample'; import { AudioSample, clampCropRectangle, CropRectangle, validateCropRectangle, VideoSample } from './sample';
@@ -377,9 +377,9 @@ export abstract class SampleCursor<
decodedTimestamps: number[] = []; decodedTimestamps: number[] = [];
maxDecodedSequenceNumber = -1; maxDecodedSequenceNumber = -1;
pumpTarget: EncodedPacket | null = null; pumpTarget: EncodedPacket | null = null;
pumpMutex = new AsyncMutex3(); pumpMutex = new AsyncMutex4();
_closed = false; _closed = false;
otherMutex = new AsyncMutex3(); // TODO: THIS IS STILL A BUGGED MUTEX! ASYNC MUTEX 4 otherMutex = new AsyncMutex4();
error: unknown = null; error: unknown = null;
errorSet = false; errorSet = false;
@@ -419,9 +419,8 @@ export abstract class SampleCursor<
return this.onDecoderError(new Error('Fake decoder error!')); return this.onDecoderError(new Error('Fake decoder error!'));
} }
const mutexPromise = this.otherMutex.request(); using lock = this.otherMutex.lock();
if (mutexPromise) await mutexPromise; if (lock.pending) await lock.ready;
using _ = this.otherMutex.lock();
while (this.decodedTimestamps.length > 0 && this.decodedTimestamps[0]! <= sample.timestamp) { while (this.decodedTimestamps.length > 0 && this.decodedTimestamps[0]! <= sample.timestamp) {
this.decodedTimestamps.shift(); this.decodedTimestamps.shift();
@@ -474,9 +473,8 @@ export abstract class SampleCursor<
} }
async onDecoderError(error: unknown) { async onDecoderError(error: unknown) {
const mutexPromise = this.otherMutex.request(); using lock = this.otherMutex.lock();
if (mutexPromise) await mutexPromise; if (lock.pending) await lock.ready;
using _ = this.otherMutex.lock();
await this.closeWithError(error); await this.closeWithError(error);
} }
@@ -519,9 +517,8 @@ export abstract class SampleCursor<
this.predictedRequests++; this.predictedRequests++;
if (!lock) { if (!lock) {
const mutexPromise = this.pumpMutex.request();
if (mutexPromise) await mutexPromise;
lock = this.pumpMutex.lock(); lock = this.pumpMutex.lock();
if (lock.pending) await lock.ready;
} }
this._ensureNotClosed(); this._ensureNotClosed();
@@ -698,9 +695,8 @@ export abstract class SampleCursor<
} }
async _nextInternal(res: ResultValue<TransformedSample | null>): Promise<Yo> { async _nextInternal(res: ResultValue<TransformedSample | null>): Promise<Yo> {
const mutexPromise = this.pumpMutex.request();
if (mutexPromise) await mutexPromise;
using lock = this.pumpMutex.lock(); using lock = this.pumpMutex.lock();
if (lock.pending) await lock.ready;
this._ensureNotClosed(); this._ensureNotClosed();
@@ -850,9 +846,8 @@ export abstract class SampleCursor<
this.predictedRequests++; this.predictedRequests++;
this.packetReader.track.input._openSampleCursors.delete(this); this.packetReader.track.input._openSampleCursors.delete(this);
const mutexPromise = this.pumpMutex.request(); using lock = this.pumpMutex.lock();
if (mutexPromise) await mutexPromise; if (lock.pending) await lock.ready;
using _ = this.pumpMutex.lock();
this._closed = true; this._closed = true;
+30 -29
View File
@@ -896,36 +896,39 @@ export class AsyncMutex2 {
} }
} }
export type AsyncMutexLock = { export interface AsyncMutexLock extends Disposable {
release: () => void; readonly pending: boolean;
[Symbol.dispose]: () => void; readonly ready: Promise<void> | null;
}; release(): void;
}
export class AsyncMutex3 { export class AsyncMutex4 {
locked = false; private locked = false;
resolverQueue: (() => void)[] = []; private resolverQueue: (() => void)[] = [];
lock(): AsyncMutexLock { lock() {
if (this.locked) { if (!this.locked) {
throw new Error('Mutex already locked.'); // Fast path
this.locked = true;
return this.createLock(false, null);
} }
this.locked = true; const { promise, resolve } = promiseWithResolvers();
this.resolverQueue.push(resolve);
return this.createLock(true, promise);
}
private createLock(pending: boolean, ready: Promise<void> | null): AsyncMutexLock {
let released = false; let released = false;
return { return {
pending,
ready,
release: () => { release: () => {
if (released) { if (released) return;
return;
}
released = true; released = true;
this.dispatch();
this.locked = false;
if (this.resolverQueue.length > 0) {
const resolve = this.resolverQueue.shift()!;
resolve();
}
}, },
[Symbol.dispose]() { [Symbol.dispose]() {
this.release(); this.release();
@@ -933,15 +936,13 @@ export class AsyncMutex3 {
}; };
} }
request() { private dispatch() {
if (!this.locked) { if (this.resolverQueue.length > 0) {
return null; const resolve = this.resolverQueue.shift()!;
resolve();
} else {
this.locked = false;
} }
const { promise, resolve } = promiseWithResolvers();
this.resolverQueue.push(resolve);
return promise;
} }
} }