Add VideoSample and AudioSample abstractions

This commit is contained in:
David Payr
2025-03-13 21:46:39 +01:00
parent d833b5ef31
commit 6eb5c02591
12 changed files with 1571 additions and 873 deletions
+160 -245
View File
@@ -17,6 +17,7 @@ import {
} from './misc';
import { EncodedPacket } from './packet';
import { fromAlaw, fromUlaw } from './pcm';
import { AudioSample, VideoSample } from './sample';
/** @public */
export type PacketRetrievalOptions = {
@@ -180,18 +181,11 @@ export class EncodedPacketSink {
}
}
export type WrappedMediaFrame<T extends VideoFrame | AudioData> = {
frame: T;
timestamp: number;
duration: number;
};
abstract class DecoderWrapper<
MediaFrame extends VideoFrame | AudioData,
WrappedFrame extends WrappedMediaFrame<MediaFrame> = WrappedMediaFrame<MediaFrame>,
MediaSample extends VideoSample | AudioSample,
> {
constructor(
public onFrame: (frame: WrappedFrame) => unknown,
public onSample: (sample: MediaSample) => unknown,
public onError: (error: DOMException) => unknown,
) {}
@@ -202,36 +196,29 @@ abstract class DecoderWrapper<
}
/** @public */
export abstract class BaseMediaFrameSink<
MediaFrame extends VideoFrame | AudioData,
/** @internal */
WrappedFrame extends WrappedMediaFrame<MediaFrame> = WrappedMediaFrame<MediaFrame>,
export abstract class BaseMediaSampleSink<
MediaSample extends VideoSample | AudioSample,
> {
/** @internal */
abstract _createDecoder(
onFrame: (frame: WrappedFrame) => unknown,
onSample: (sample: MediaSample) => unknown,
onError: (error: DOMException) => unknown
): Promise<DecoderWrapper<MediaFrame>>;
): Promise<DecoderWrapper<MediaSample>>;
/** @internal */
abstract _createPacketSink(): EncodedPacketSink;
/** @internal */
private _duplicateFrame(frame: WrappedFrame) {
return structuredClone(frame);
}
/** @internal */
protected mediaFramesInRange(
protected mediaSamplesInRange(
startTimestamp = 0,
endTimestamp = Infinity,
): AsyncGenerator<WrappedFrame, void, unknown> {
): AsyncGenerator<MediaSample, void, unknown> {
validateTimestamp(startTimestamp);
validateTimestamp(endTimestamp);
const MAX_QUEUE_SIZE = 8;
const frameQueue: WrappedFrame[] = [];
let firstFrameQueued = false;
let lastFrame: WrappedFrame | null = null;
const sampleQueue: MediaSample[] = [];
let firstSampleQueued = false;
let lastSample: MediaSample | null = null;
let { promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers();
let { promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers();
let decoderIsFlushed = false;
@@ -246,38 +233,38 @@ export abstract class BaseMediaFrameSink<
// The following is the "pump" process that keeps pumping packets into the decoder
(async () => {
const decoderError = new Error();
const decoder = await this._createDecoder((wrappedFrame) => {
const decoder = await this._createDecoder((sample) => {
onQueueDequeue();
if (wrappedFrame.timestamp >= endTimestamp) {
if (sample.timestamp >= endTimestamp) {
ended = true;
}
if (ended) {
wrappedFrame.frame.close();
sample.close();
return;
}
if (lastFrame) {
if (wrappedFrame.timestamp > startTimestamp) {
if (lastSample) {
if (sample.timestamp > startTimestamp) {
// We don't know ahead of time what the first first is. This is because the first first is the
// last first whose timestamp is less than or equal to the start timestamp. Therefore we need to
// wait for the first first after the start timestamp, and then we'll know that the previous
// first was the first first.
frameQueue.push(lastFrame);
firstFrameQueued = true;
sampleQueue.push(lastSample);
firstSampleQueued = true;
} else {
lastFrame.frame.close();
lastSample.close();
}
}
if (wrappedFrame.timestamp >= startTimestamp) {
frameQueue.push(wrappedFrame);
firstFrameQueued = true;
if (sample.timestamp >= startTimestamp) {
sampleQueue.push(sample);
firstSampleQueued = true;
}
lastFrame = firstFrameQueued ? null : wrappedFrame;
lastSample = firstSampleQueued ? null : sample;
if (frameQueue.length > 0) {
if (sampleQueue.length > 0) {
onQueueNotEmpty();
({ promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers());
}
@@ -319,7 +306,7 @@ export abstract class BaseMediaFrameSink<
await packets.next(); // Skip the start packet as we already have it
while (currentPacket && !ended) {
if (frameQueue.length + decoder.getDecodeQueueSize() > MAX_QUEUE_SIZE) {
if (sampleQueue.length + decoder.getDecodeQueueSize() > MAX_QUEUE_SIZE) {
({ promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers());
await queueDequeue;
continue;
@@ -340,8 +327,8 @@ export abstract class BaseMediaFrameSink<
if (!terminated) await decoder.flush();
decoder.close();
if (!firstFrameQueued && lastFrame) {
frameQueue.push(lastFrame);
if (!firstSampleQueued && lastSample) {
sampleQueue.push(lastSample);
}
decoderIsFlushed = true;
@@ -360,8 +347,8 @@ export abstract class BaseMediaFrameSink<
return { value: undefined, done: true };
} else if (outOfBandError) {
throw outOfBandError;
} else if (frameQueue.length > 0) {
const value = frameQueue.shift()!;
} else if (sampleQueue.length > 0) {
const value = sampleQueue.shift()!;
onQueueDequeue();
return { value, done: false };
} else if (!decoderIsFlushed) {
@@ -377,10 +364,10 @@ export abstract class BaseMediaFrameSink<
onQueueDequeue();
onQueueNotEmpty();
lastFrame?.frame.close();
lastSample?.close();
for (const frame of frameQueue) {
frame.frame.close();
for (const sample of sampleQueue) {
sample.close();
}
return { value: undefined, done: true };
@@ -395,15 +382,15 @@ export abstract class BaseMediaFrameSink<
}
/** @internal */
protected mediaFramesAtTimestamps(
protected mediaSamplesAtTimestamps(
timestamps: AnyIterable<number>,
): AsyncGenerator<WrappedFrame | null, void, unknown> {
): AsyncGenerator<MediaSample | null, void, unknown> {
validateAnyIterable(timestamps);
const timestampIterator = toAsyncIterator(timestamps);
const timestampsOfInterest: number[] = [];
const MAX_QUEUE_SIZE = 8;
const frameQueue: (WrappedFrame | null)[] = [];
const sampleQueue: (MediaSample | null)[] = [];
let { promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers();
let { promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers();
let decoderIsFlushed = false;
@@ -414,9 +401,9 @@ export abstract class BaseMediaFrameSink<
// the consumer.
let outOfBandError = null as Error | null;
let lastUsedFrame = null as WrappedFrame | null;
const pushToQueue = (frame: WrappedFrame | null) => {
frameQueue.push(frame);
let lastUsedSample = null as MediaSample | null;
const pushToQueue = (sample: MediaSample | null) => {
sampleQueue.push(sample);
onQueueNotEmpty();
({ promise: queueNotEmpty, resolve: onQueueNotEmpty } = promiseWithResolvers());
};
@@ -424,29 +411,29 @@ export abstract class BaseMediaFrameSink<
// The following is the "pump" process that keeps pumping packets into the decoder
(async () => {
const decoderError = new Error();
const decoder = await this._createDecoder((wrappedFrame) => {
const decoder = await this._createDecoder((sample) => {
onQueueDequeue();
if (terminated) {
wrappedFrame.frame.close();
sample.close();
return;
}
let frameUsed = false;
let sampleUsed = false;
while (
timestampsOfInterest.length > 0
&& wrappedFrame.timestamp - timestampsOfInterest[0]! > -1e-10 // Give it a little epsilon
&& sample.timestamp - timestampsOfInterest[0]! > -1e-10 // Give it a little epsilon
) {
pushToQueue(this._duplicateFrame(wrappedFrame));
pushToQueue(sample.clone() as MediaSample);
timestampsOfInterest.shift();
frameUsed = true;
sampleUsed = true;
}
if (frameUsed) {
lastUsedFrame?.frame.close();
lastUsedFrame = wrappedFrame;
if (sampleUsed) {
lastUsedSample?.close();
lastUsedSample = sample;
} else {
wrappedFrame.frame.close();
sample.close();
}
}, (error) => {
if (!outOfBandError) {
@@ -463,7 +450,7 @@ export abstract class BaseMediaFrameSink<
for await (const timestamp of timestampIterator) {
validateTimestamp(timestamp);
while (frameQueue.length + decoder.getDecodeQueueSize() > MAX_QUEUE_SIZE && !terminated) {
while (sampleQueue.length + decoder.getDecodeQueueSize() > MAX_QUEUE_SIZE && !terminated) {
({ promise: queueDequeue, resolve: onQueueDequeue } = promiseWithResolvers());
await queueDequeue;
}
@@ -501,10 +488,10 @@ export abstract class BaseMediaFrameSink<
targetPacket.sequenceNumber === lastPacket.sequenceNumber
&& timestampsOfInterest.length === 0
) {
// Special case: We have a repeat packet, but the frame for that packet has already been
// decoded. Therefore, we need to push the frame here instead of in the decoder callback.
if (lastUsedFrame) {
pushToQueue(this._duplicateFrame(lastUsedFrame));
// Special case: We have a repeat packet, but the sample for that packet has already been
// decoded. Therefore, we need to push the sample here instead of in the decoder callback.
if (lastUsedSample) {
pushToQueue(lastUsedSample.clone() as MediaSample);
}
} else {
timestampsOfInterest.push(targetPacket.timestamp);
@@ -528,7 +515,7 @@ export abstract class BaseMediaFrameSink<
if (!terminated) {
await decoder.flush();
lastUsedFrame?.frame.close();
lastUsedSample?.close();
}
decoder.close();
@@ -548,8 +535,8 @@ export abstract class BaseMediaFrameSink<
return { value: undefined, done: true };
} else if (outOfBandError) {
throw outOfBandError;
} else if (frameQueue.length > 0) {
const value = frameQueue.shift();
} else if (sampleQueue.length > 0) {
const value = sampleQueue.shift();
assert(value !== undefined);
onQueueDequeue();
return { value, done: false };
@@ -565,10 +552,10 @@ export abstract class BaseMediaFrameSink<
onQueueDequeue();
onQueueNotEmpty();
for (const frame of frameQueue) {
frame?.frame.close();
for (const sample of sampleQueue) {
sample?.close();
}
lastUsedFrame?.frame.close();
lastUsedSample?.close();
return { value: undefined, done: true };
},
@@ -582,44 +569,45 @@ export abstract class BaseMediaFrameSink<
}
}
class VideoDecoderWrapper extends DecoderWrapper<VideoFrame> {
class VideoDecoderWrapper extends DecoderWrapper<VideoSample> {
decoder: VideoDecoder | null = null;
customDecoder: CustomVideoDecoder | null = null;
lastCustomDecoderPromise = Promise.resolve();
customDecoderQueueSize = 0;
frameQueue: VideoFrame[] = [];
sampleQueue: VideoSample[] = [];
constructor(
onFrame: (frame: WrappedMediaFrame<VideoFrame>) => unknown,
onSample: (sample: VideoSample) => unknown,
onError: (error: DOMException) => unknown,
codec: VideoCodec,
decoderConfig: VideoDecoderConfig,
public rotation: Rotation,
public timeResolution: number,
) {
super(onFrame, onError);
super(onSample, onError);
const frameHandler = (frame: VideoFrame) => {
const sampleHandler = (sample: VideoSample) => {
// For correct B-frame handling, we don't just hand over the frames directly but instead add them to a
// queue, because we want to ensure frames are emitted in presentation order. We flush the queue each time
// we receive a frame with a timestamp larger than the highest we've seen so far, as we can sure that is
// not a B-frame. Typically, WebCodecs automatically guarantees that frames are emitted in presentation
// order, but some browsers (Safari) don't always follow this rule.
if (this.frameQueue.length > 0 && (frame.timestamp >= last(this.frameQueue)!.timestamp)) {
for (const frame of this.frameQueue) {
this.wrapAndEmitFrame(frame);
if (this.sampleQueue.length > 0 && (sample.timestamp >= last(this.sampleQueue)!.timestamp)) {
for (const sample of this.sampleQueue) {
this.finalizeAndEmitSample(sample);
}
this.frameQueue.length = 0;
this.sampleQueue.length = 0;
}
const insertionIndex = binarySearchLessOrEqual(
this.frameQueue,
frame.timestamp,
this.sampleQueue,
sample.timestamp,
x => x.timestamp,
);
this.frameQueue.splice(insertionIndex + 1, 0, frame);
this.sampleQueue.splice(insertionIndex + 1, 0, sample);
};
const MatchingCustomDecoder = customVideoDecoders.find(x => x.supports(codec, decoderConfig));
@@ -628,28 +616,24 @@ class VideoDecoderWrapper extends DecoderWrapper<VideoFrame> {
this.customDecoder = new MatchingCustomDecoder() as CustomVideoDecoder;
this.customDecoder.codec = codec;
this.customDecoder.config = decoderConfig;
this.customDecoder.onFrame = frameHandler;
this.customDecoder.onSample = sampleHandler;
this.customDecoder.init();
} else {
this.decoder = new VideoDecoder({
output: frameHandler,
output: frame => sampleHandler(new VideoSample(frame)),
error: onError,
});
this.decoder.configure(decoderConfig);
}
}
wrapAndEmitFrame(frame: VideoFrame) {
// Round the microsecond timestamps to the time resolution
const timestamp = Math.round(frame.timestamp / 1e6 * this.timeResolution) / this.timeResolution;
const duration = Math.round((frame.duration ?? 0) / 1e6 * this.timeResolution) / this.timeResolution;
finalizeAndEmitSample(sample: VideoSample) {
// Round the timestamps to the time resolution
sample.setTimestamp(Math.round(sample.timestamp * this.timeResolution) / this.timeResolution);
sample.setDuration(Math.round(sample.duration * this.timeResolution) / this.timeResolution);
this.onFrame({
frame,
timestamp,
duration,
});
this.onSample(sample);
}
getDecodeQueueSize() {
@@ -683,10 +667,10 @@ class VideoDecoderWrapper extends DecoderWrapper<VideoFrame> {
await this.decoder.flush();
}
for (const frame of this.frameQueue) {
this.wrapAndEmitFrame(frame);
for (const sample of this.sampleQueue) {
this.finalizeAndEmitSample(sample);
}
this.frameQueue.length = 0;
this.sampleQueue.length = 0;
}
close() {
@@ -697,22 +681,15 @@ class VideoDecoderWrapper extends DecoderWrapper<VideoFrame> {
this.decoder.close();
}
for (const frame of this.frameQueue) {
frame.close();
for (const sample of this.sampleQueue) {
sample.close();
}
this.frameQueue.length = 0;
this.sampleQueue.length = 0;
}
}
/** @public */
export type WrappedVideoFrame = {
frame: VideoFrame;
timestamp: number;
duration: number;
};
/** @public */
export class VideoFrameSink extends BaseMediaFrameSink<VideoFrame> {
export class VideoSampleSink extends BaseMediaSampleSink<VideoSample> {
/** @internal */
_videoTrack: InputVideoTrack;
@@ -728,7 +705,7 @@ export class VideoFrameSink extends BaseMediaFrameSink<VideoFrame> {
/** @internal */
async _createDecoder(
onFrame: (frame: WrappedMediaFrame<VideoFrame>) => unknown,
onSample: (sample: VideoSample) => unknown,
onError: (error: DOMException) => unknown,
) {
if (!(await this._videoTrack.canDecode())) {
@@ -739,11 +716,12 @@ export class VideoFrameSink extends BaseMediaFrameSink<VideoFrame> {
}
const codec = this._videoTrack.codec;
const rotation = this._videoTrack.rotation;
const decoderConfig = await this._videoTrack.getDecoderConfig();
const timeResolution = this._videoTrack.timeResolution;
assert(codec && decoderConfig);
return new VideoDecoderWrapper(onFrame, onError, codec, decoderConfig, timeResolution);
return new VideoDecoderWrapper(onSample, onError, codec, decoderConfig, rotation, timeResolution);
}
/** @internal */
@@ -751,36 +729,21 @@ export class VideoFrameSink extends BaseMediaFrameSink<VideoFrame> {
return new EncodedPacketSink(this._videoTrack);
}
/** @internal */
_wrappedFrameToWrappedVideoFrame(frame: WrappedMediaFrame<VideoFrame>): WrappedVideoFrame {
return {
frame: frame.frame,
timestamp: frame.timestamp,
duration: frame.duration,
};
}
async getFrame(timestamp: number) {
async getSample(timestamp: number) {
validateTimestamp(timestamp);
for await (const frame of this.mediaFramesAtTimestamps([timestamp])) {
return frame && this._wrappedFrameToWrappedVideoFrame(frame);
for await (const sample of this.mediaSamplesAtTimestamps([timestamp])) {
return sample;
}
throw new Error('Internal error: Iterator returned nothing.');
}
frames(startTimestamp = 0, endTimestamp = Infinity) {
return mapAsyncGenerator(
this.mediaFramesInRange(startTimestamp, endTimestamp),
frame => this._wrappedFrameToWrappedVideoFrame(frame),
);
samples(startTimestamp = 0, endTimestamp = Infinity) {
return this.mediaSamplesInRange(startTimestamp, endTimestamp);
}
framesAtTimestamps(timestamps: AnyIterable<number>) {
return mapAsyncGenerator(
this.mediaFramesAtTimestamps(timestamps),
frame => frame && this._wrappedFrameToWrappedVideoFrame(frame),
);
samplesAtTimestamps(timestamps: AnyIterable<number>) {
return this.mediaSamplesAtTimestamps(timestamps);
}
}
@@ -813,7 +776,7 @@ export class CanvasSink {
/** @internal */
_rotation: Rotation;
/** @internal */
_videoFrameSink: VideoFrameSink;
_videoSampleSink: VideoSampleSink;
/** @internal */
_canvasPool: (HTMLCanvasElement | OffscreenCanvas | null)[];
/** @internal */
@@ -877,12 +840,12 @@ export class CanvasSink {
this._height = height;
this._rotation = rotation;
this._fit = options.fit ?? 'fill';
this._videoFrameSink = new VideoFrameSink(videoTrack);
this._videoSampleSink = new VideoSampleSink(videoTrack);
this._canvasPool = Array.from({ length: options.poolSize ?? 0 }, () => null);
}
/** @internal */
_videoFrameToWrappedCanvas(frame: WrappedVideoFrame): WrappedCanvas {
_videoSampleToWrappedCanvas(sample: VideoSample): WrappedCanvas {
let canvas = this._canvasPool[this._nextCanvasIndex];
if (!canvas) {
if (typeof OffscreenCanvas !== 'undefined') {
@@ -903,12 +866,13 @@ export class CanvasSink {
this._nextCanvasIndex = (this._nextCanvasIndex + 1) % this._canvasPool.length;
}
const context = canvas.getContext('2d', { alpha: false });
const context
= canvas.getContext('2d', { alpha: false }) as CanvasRenderingContext2D | OffscreenCanvasRenderingContext2D;
assert(context);
context.resetTransform();
// These variables specify where the final frame will be drawn on the canvas
// These variables specify where the final sample will be drawn on the canvas
let dx: number;
let dy: number;
let newWidth: number;
@@ -920,15 +884,15 @@ export class CanvasSink {
newWidth = this._width;
newHeight = this._height;
} else {
const [frameWidth, frameHeight] = this._rotation % 180 === 0
? [frame.frame.codedWidth, frame.frame.codedHeight]
: [frame.frame.codedHeight, frame.frame.codedWidth];
const [sampleWidth, sampleHeight] = this._rotation % 180 === 0
? [sample.codedWidth, sample.codedHeight]
: [sample.codedHeight, sample.codedWidth];
const scale = this._fit === 'contain'
? Math.min(this._width / frameWidth, this._height / frameHeight)
: Math.max(this._width / frameWidth, this._height / frameHeight);
newWidth = frameWidth * scale;
newHeight = frameHeight * scale;
? Math.min(this._width / sampleWidth, this._height / sampleHeight)
: Math.max(this._width / sampleWidth, this._height / sampleHeight);
newWidth = sampleWidth * scale;
newHeight = sampleHeight * scale;
dx = (this._width - newWidth) / 2;
dy = (this._height - newHeight) / 2;
}
@@ -936,46 +900,46 @@ export class CanvasSink {
const aspectRatioChange = this._rotation % 180 === 0 ? 1 : newWidth / newHeight;
context.translate(this._width / 2, this._height / 2);
context.rotate(this._rotation * Math.PI / 180);
// This aspect ratio compensation is done so that we can draw the frame with the intended dimensions and
// This aspect ratio compensation is done so that we can draw the sample with the intended dimensions and
// don't need to think about how those dimensions change after the rotation
context.scale(1 / aspectRatioChange, aspectRatioChange);
context.translate(-this._width / 2, -this._height / 2);
context.drawImage(frame.frame, dx, dy, newWidth, newHeight);
context.drawImage(sample.toCanvasImageSource(), dx, dy, newWidth, newHeight);
const result = {
canvas,
timestamp: frame.timestamp,
duration: frame.duration,
timestamp: sample.timestamp,
duration: sample.duration,
};
frame.frame.close();
sample.close();
return result;
}
async getCanvas(timestamp: number) {
validateTimestamp(timestamp);
const frame = await this._videoFrameSink.getFrame(timestamp);
return frame && this._videoFrameToWrappedCanvas(frame);
const sample = await this._videoSampleSink.getSample(timestamp);
return sample && this._videoSampleToWrappedCanvas(sample);
}
canvases(startTimestamp = 0, endTimestamp = Infinity) {
return mapAsyncGenerator(
this._videoFrameSink.frames(startTimestamp, endTimestamp),
frame => this._videoFrameToWrappedCanvas(frame),
this._videoSampleSink.samples(startTimestamp, endTimestamp),
sample => this._videoSampleToWrappedCanvas(sample),
);
}
canvasesAtTimestamps(timestamps: AnyIterable<number>) {
return mapAsyncGenerator(
this._videoFrameSink.framesAtTimestamps(timestamps),
frame => frame && this._videoFrameToWrappedCanvas(frame),
this._videoSampleSink.samplesAtTimestamps(timestamps),
sample => sample && this._videoSampleToWrappedCanvas(sample),
);
}
}
class AudioDecoderWrapper extends DecoderWrapper<AudioData> {
class AudioDecoderWrapper extends DecoderWrapper<AudioSample> {
decoder: AudioDecoder | null = null;
customDecoder: CustomAudioDecoder | null = null;
@@ -983,25 +947,20 @@ class AudioDecoderWrapper extends DecoderWrapper<AudioData> {
customDecoderQueueSize = 0;
constructor(
onData: (data: WrappedMediaFrame<AudioData>) => unknown,
onSample: (sample: AudioSample) => unknown,
onError: (error: DOMException) => unknown,
codec: AudioCodec,
decoderConfig: AudioDecoderConfig,
) {
super(onData, onError);
super(onSample, onError);
const dataHandler = (data: AudioData) => {
const sampleHandler = (sample: AudioSample) => {
const sampleRate = decoderConfig.sampleRate;
// Round the microsecond timestamps to the sample rate
const timestamp = Math.round(data.timestamp / 1e6 * sampleRate) / sampleRate;
const duration = Math.round(data.duration / 1e6 * sampleRate) / sampleRate;
// Round the timestamp to the sample rate
sample.setTimestamp(Math.round(sample.timestamp * sampleRate) / sampleRate);
onData({
frame: data,
timestamp,
duration,
});
onSample(sample);
};
const MatchingCustomDecoder = customAudioDecoders.find(x => x.supports(codec, decoderConfig));
@@ -1010,12 +969,12 @@ class AudioDecoderWrapper extends DecoderWrapper<AudioData> {
this.customDecoder = new MatchingCustomDecoder() as CustomAudioDecoder;
this.customDecoder.codec = codec;
this.customDecoder.config = decoderConfig;
this.customDecoder.onData = dataHandler;
this.customDecoder.onSample = sampleHandler;
this.customDecoder.init();
} else {
this.decoder = new AudioDecoder({
output: dataHandler,
output: data => sampleHandler(new AudioSample(data)),
error: onError,
});
this.decoder.configure(decoderConfig);
@@ -1066,7 +1025,7 @@ class AudioDecoderWrapper extends DecoderWrapper<AudioData> {
// There are a lot of PCM variants not natively supported by the browser and by AudioData. Therefore we need a simple
// decoder that maps any input PCM format into a PCM format supported by the browser.
class PcmAudioDecoderWrapper extends DecoderWrapper<AudioData> {
class PcmAudioDecoderWrapper extends DecoderWrapper<AudioSample> {
codec: PcmAudioCodec;
inputSampleSize: 1 | 2 | 3 | 4;
@@ -1081,11 +1040,11 @@ class PcmAudioDecoderWrapper extends DecoderWrapper<AudioData> {
currentTimestamp: number | null = null;
constructor(
onData: (data: WrappedMediaFrame<AudioData>) => unknown,
onSample: (sample: AudioSample) => unknown,
onError: (error: DOMException) => unknown,
public decoderConfig: AudioDecoderConfig,
) {
super(onData, onError);
super(onSample, onError);
assert((PCM_AUDIO_CODECS as readonly string[]).includes(decoderConfig.codec));
this.codec = decoderConfig.codec as PcmAudioCodec;
@@ -1207,21 +1166,16 @@ class PcmAudioDecoderWrapper extends DecoderWrapper<AudioData> {
const preciseTimestamp = this.currentTimestamp;
this.currentTimestamp += preciseDuration;
const audioData = new AudioData({
const audioSample = new AudioSample({
format: this.outputFormat,
data: outputBuffer,
numberOfChannels: this.decoderConfig.numberOfChannels,
sampleRate: this.decoderConfig.sampleRate,
numberOfFrames,
timestamp: 1e6 * preciseTimestamp,
timestamp: preciseTimestamp,
});
// Since all other decoders are async, we'll make this one behave async as well
queueMicrotask(() => this.onFrame({
frame: audioData,
timestamp: preciseTimestamp,
duration: preciseDuration,
}));
this.onSample(audioSample);
}
async flush() {
@@ -1234,14 +1188,7 @@ class PcmAudioDecoderWrapper extends DecoderWrapper<AudioData> {
}
/** @public */
export type WrappedAudioData = {
data: AudioData;
timestamp: number;
duration: number;
};
/** @public */
export class AudioDataSink extends BaseMediaFrameSink<AudioData> {
export class AudioSampleSink extends BaseMediaSampleSink<AudioSample> {
/** @internal */
_audioTrack: InputAudioTrack;
@@ -1257,7 +1204,7 @@ export class AudioDataSink extends BaseMediaFrameSink<AudioData> {
/** @internal */
async _createDecoder(
onData: (data: WrappedMediaFrame<AudioData>) => unknown,
onSample: (sample: AudioSample) => unknown,
onError: (error: DOMException) => unknown,
) {
if (!(await this._audioTrack.canDecode())) {
@@ -1272,47 +1219,32 @@ export class AudioDataSink extends BaseMediaFrameSink<AudioData> {
assert(codec && decoderConfig);
if ((PCM_AUDIO_CODECS as readonly string[]).includes(decoderConfig.codec)) {
return new PcmAudioDecoderWrapper(onData, onError, decoderConfig);
return new PcmAudioDecoderWrapper(onSample, onError, decoderConfig);
} else {
return new AudioDecoderWrapper(onData, onError, codec, decoderConfig);
return new AudioDecoderWrapper(onSample, onError, codec, decoderConfig);
}
}
/** @internal */
_wrappedFrameToWrappedAudioData(frame: WrappedMediaFrame<AudioData>): WrappedAudioData {
return {
data: frame.frame,
timestamp: frame.timestamp,
duration: frame.duration,
};
}
/** @internal */
_createPacketSink() {
return new EncodedPacketSink(this._audioTrack);
}
async getData(timestamp: number) {
async getSample(timestamp: number) {
validateTimestamp(timestamp);
for await (const data of this.mediaFramesAtTimestamps([timestamp])) {
return data && this._wrappedFrameToWrappedAudioData(data);
for await (const sample of this.mediaSamplesAtTimestamps([timestamp])) {
return sample;
}
throw new Error('Internal error: Iterator returned nothing.');
}
data(startTimestamp = 0, endTimestamp = Infinity) {
return mapAsyncGenerator(
this.mediaFramesInRange(startTimestamp, endTimestamp),
data => this._wrappedFrameToWrappedAudioData(data),
);
samples(startTimestamp = 0, endTimestamp = Infinity) {
return this.mediaSamplesInRange(startTimestamp, endTimestamp);
}
dataAtTimestamps(timestamps: AnyIterable<number>) {
return mapAsyncGenerator(
this.mediaFramesAtTimestamps(timestamps),
data => data && this._wrappedFrameToWrappedAudioData(data),
);
samplesAtTimestamps(timestamps: AnyIterable<number>) {
return this.mediaSamplesAtTimestamps(timestamps);
}
}
@@ -1326,60 +1258,43 @@ export type WrappedAudioBuffer = {
/** @public */
export class AudioBufferSink {
/** @internal */
_audioDataSink: AudioDataSink;
_audioSampleSink: AudioSampleSink;
constructor(audioTrack: InputAudioTrack) {
if (!(audioTrack instanceof InputAudioTrack)) {
throw new TypeError('audioTrack must be an InputAudioTrack.');
}
this._audioDataSink = new AudioDataSink(audioTrack);
this._audioSampleSink = new AudioSampleSink(audioTrack);
}
/** @internal */
_audioDataToWrappedArrayBuffer(data: WrappedAudioData): WrappedAudioBuffer {
const audioBuffer = new AudioBuffer({
numberOfChannels: data.data.numberOfChannels,
length: data.data.numberOfFrames,
sampleRate: data.data.sampleRate,
});
// All user agents are required to support conversion to f32-planar
const dataBytes = new Float32Array(data.data.allocationSize({ planeIndex: 0, format: 'f32-planar' }) / 4);
for (let i = 0; i < data.data.numberOfChannels; i++) {
data.data.copyTo(dataBytes, { planeIndex: i, format: 'f32-planar' });
audioBuffer.copyToChannel(dataBytes, i);
}
const result = {
buffer: audioBuffer,
timestamp: data.timestamp,
duration: data.duration,
_audioSampleToWrappedArrayBuffer(sample: AudioSample): WrappedAudioBuffer {
return {
buffer: sample.toAudioBuffer(),
timestamp: sample.timestamp,
duration: sample.duration,
};
data.data.close();
return result;
}
async getBuffer(timestamp: number) {
validateTimestamp(timestamp);
const data = await this._audioDataSink.getData(timestamp);
return data && this._audioDataToWrappedArrayBuffer(data);
const data = await this._audioSampleSink.getSample(timestamp);
return data && this._audioSampleToWrappedArrayBuffer(data);
}
buffers(startTimestamp = 0, endTimestamp = Infinity) {
return mapAsyncGenerator(
this._audioDataSink.data(startTimestamp, endTimestamp),
data => this._audioDataToWrappedArrayBuffer(data),
this._audioSampleSink.samples(startTimestamp, endTimestamp),
data => this._audioSampleToWrappedArrayBuffer(data),
);
}
buffersAtTimestamps(timestamps: AnyIterable<number>) {
return mapAsyncGenerator(
this._audioDataSink.dataAtTimestamps(timestamps),
data => data && this._audioDataToWrappedArrayBuffer(data),
this._audioSampleSink.samplesAtTimestamps(timestamps),
data => data && this._audioSampleToWrappedArrayBuffer(data),
);
}
}