Big realtime playback refactor

- Added Worker- or AudioContext-based fallbacks in case MediaStreamTrackProcessor is not available
- Fixed fMP4-streamed files not playing in Safari
- Added better error handling for MediaStream sources
- Improved the live recording example to be more robust
This commit is contained in:
Vanilagy
2025-07-27 12:13:03 +02:00
parent 9cc38329f2
commit d38ad22559
12 changed files with 690 additions and 268 deletions
+545 -235
View File
@@ -24,7 +24,7 @@ import {
VideoCodec,
} from './codec';
import { OutputAudioTrack, OutputSubtitleTrack, OutputTrack, OutputVideoTrack } from './output';
import { assert, assertNever, CallSerializer, clamp, setInt24, setUint24 } from './misc';
import { assert, assertNever, CallSerializer, clamp, promiseWithResolvers, setInt24, setUint24 } from './misc';
import { Muxer } from './muxer';
import { SubtitleParser } from './subtitles';
import { toAlaw, toUlaw } from './pcm';
@@ -78,7 +78,7 @@ export abstract class MediaSource {
}
/** @internal */
_start() {}
async _start() {}
/** @internal */
async _flushAndClose() {}
@@ -272,86 +272,94 @@ class VideoEncoderWrapper {
constructor(private source: VideoSource, private encodingConfig: VideoEncodingConfig) {}
async add(videoSample: VideoSample, shouldClose: boolean, encodeOptions?: VideoEncoderEncodeOptions) {
this.checkForEncoderError();
this.source._ensureValidAdd();
try {
this.checkForEncoderError();
this.source._ensureValidAdd();
// Ensure video sample size remains constant
if (this.lastWidth !== null && this.lastHeight !== null) {
if (videoSample.codedWidth !== this.lastWidth || videoSample.codedHeight !== this.lastHeight) {
throw new Error(
`Video sample size must remain constant. Expected ${this.lastWidth}x${this.lastHeight},`
+ ` got ${videoSample.codedWidth}x${videoSample.codedHeight}.`,
);
}
} else {
this.lastWidth = videoSample.codedWidth;
this.lastHeight = videoSample.codedHeight;
}
if (!this.encoderInitialized) {
if (!this.ensureEncoderPromise) {
void this.ensureEncoder(videoSample);
// Ensure video sample size remains constant
if (this.lastWidth !== null && this.lastHeight !== null) {
if (videoSample.codedWidth !== this.lastWidth || videoSample.codedHeight !== this.lastHeight) {
throw new Error(
`Video sample size must remain constant. Expected ${this.lastWidth}x${this.lastHeight},`
+ ` got ${videoSample.codedWidth}x${videoSample.codedHeight}.`,
);
}
} else {
this.lastWidth = videoSample.codedWidth;
this.lastHeight = videoSample.codedHeight;
}
// No, this "if" statement is not useless. Sometimes, the above call to `ensureEncoder` might have
// synchronously completed and the encoder is already initialized. In this case, we don't need to await the
// promise anymore. This also fixes nasty async race condition bugs when multiple code paths are calling
// this method: It's important that the call that initialized the encoder go through this code first.
if (!this.encoderInitialized) {
await this.ensureEncoderPromise;
if (!this.ensureEncoderPromise) {
void this.ensureEncoder(videoSample);
}
// No, this "if" statement is not useless. Sometimes, the above call to `ensureEncoder` might have
// synchronously completed and the encoder is already initialized. In this case, we don't need to await
// the promise anymore. This also fixes nasty async race condition bugs when multiple code paths are
// calling this method: It's important that the call that initialized the encoder go through this
// code first.
if (!this.encoderInitialized) {
await this.ensureEncoderPromise;
}
}
}
assert(this.encoderInitialized);
assert(this.encoderInitialized);
const keyFrameInterval = this.encodingConfig.keyFrameInterval ?? 5;
const multipleOfKeyFrameInterval = Math.floor(videoSample.timestamp / keyFrameInterval);
const keyFrameInterval = this.encodingConfig.keyFrameInterval ?? 5;
const multipleOfKeyFrameInterval = Math.floor(videoSample.timestamp / keyFrameInterval);
// Ensure a key frame every keyFrameInterval seconds. It is important that all video tracks follow the same
// "key frame" rhythm, because aligned key frames are required to start new fragments in ISOBMFF or clusters
// in Matroska (or at least desirable).
const finalEncodeOptions = {
...encodeOptions,
keyFrame: encodeOptions?.keyFrame
|| keyFrameInterval === 0
|| multipleOfKeyFrameInterval !== this.lastMultipleOfKeyFrameInterval,
};
this.lastMultipleOfKeyFrameInterval = multipleOfKeyFrameInterval;
// Ensure a key frame every keyFrameInterval seconds. It is important that all video tracks follow the same
// "key frame" rhythm, because aligned key frames are required to start new fragments in ISOBMFF or clusters
// in Matroska (or at least desirable).
const finalEncodeOptions = {
...encodeOptions,
keyFrame: encodeOptions?.keyFrame
|| keyFrameInterval === 0
|| multipleOfKeyFrameInterval !== this.lastMultipleOfKeyFrameInterval,
};
this.lastMultipleOfKeyFrameInterval = multipleOfKeyFrameInterval;
if (this.customEncoder) {
this.customEncoderQueueSize++;
const promise = this.customEncoderCallSerializer
.call(() => this.customEncoder!.encode(videoSample, finalEncodeOptions))
.then(() => {
this.customEncoderQueueSize--;
if (this.customEncoder) {
this.customEncoderQueueSize++;
const promise = this.customEncoderCallSerializer
.call(() => this.customEncoder!.encode(videoSample, finalEncodeOptions))
.then(() => {
this.customEncoderQueueSize--;
if (shouldClose) {
videoSample.close();
}
})
.catch((error: Error) => {
this.encoderError ??= error;
});
if (shouldClose) {
videoSample.close();
}
})
.catch((error: Error) => {
this.encoderError ??= error;
});
if (this.customEncoderQueueSize >= 4) {
await promise;
if (this.customEncoderQueueSize >= 4) {
await promise;
}
} else {
assert(this.encoder);
const videoFrame = videoSample.toVideoFrame();
this.encoder.encode(videoFrame, finalEncodeOptions);
videoFrame.close();
if (shouldClose) {
videoSample.close();
}
// We need to do this after sending the frame to the encoder as the frame otherwise might be closed
if (this.encoder.encodeQueueSize >= 4) {
await new Promise(resolve => this.encoder!.addEventListener('dequeue', resolve, { once: true }));
}
}
} else {
assert(this.encoder);
const videoFrame = videoSample.toVideoFrame();
this.encoder.encode(videoFrame, finalEncodeOptions);
videoFrame.close();
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
} finally {
if (shouldClose) {
// Make sure it's always closed, even if there was an error
videoSample.close();
}
// We need to do this after sending the frame to the encoder as the frame otherwise might be closed
if (this.encoder.encodeQueueSize >= 4) {
await new Promise(resolve => this.encoder!.addEventListener('dequeue', resolve, { once: true }));
}
}
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
}
private async ensureEncoder(videoSample: VideoSample) {
@@ -575,6 +583,20 @@ export class MediaStreamVideoTrackSource extends VideoSource {
private _abortController: AbortController | null = null;
/** @internal */
private _track: MediaStreamVideoTrack;
/** @internal */
private _workerTrackId: number | null = null;
/** @internal */
private _workerListener: ((event: MessageEvent) => void) | null = null;
/** @internal */
private _promiseWithResolvers = promiseWithResolvers();
/** @internal */
private _errorPromiseAccessed = false;
/** A promise that rejects upon any error within this source. This promise never resolves. */
get errorPromise() {
this._errorPromiseAccessed = true;
return this._promiseWithResolvers.promise;
}
constructor(track: MediaStreamVideoTrack, encodingConfig: VideoEncodingConfig) {
if (!(track instanceof MediaStreamTrack) || track.kind !== 'video') {
@@ -593,41 +615,102 @@ export class MediaStreamVideoTrackSource extends VideoSource {
}
/** @internal */
override _start() {
override async _start() {
if (!this._errorPromiseAccessed) {
console.warn(
'Make sure not to ignore the `errorPromise` field on MediaStreamVideoTrackSource, so that any internal'
+ ' errors get bubbled up properly.',
);
}
this._abortController = new AbortController();
let frameReceived = false;
let firstVideoFrameTimestamp: number | null = null;
let errored = false;
const processor = new MediaStreamTrackProcessor({ track: this._track });
const consumer = new WritableStream<VideoFrame>({
write: (videoFrame) => {
if (!frameReceived) {
setMediaStreamTimestampOffset(this, videoFrame);
frameReceived = true;
const onVideoFrame = (videoFrame: VideoFrame) => {
if (errored) {
videoFrame.close();
return;
}
if (firstVideoFrameTimestamp === null) {
firstVideoFrameTimestamp = videoFrame.timestamp / 1e6;
const muxer = this._connectedTrack!.output._muxer;
if (muxer.firstMediaStreamTimestamp === null) {
muxer.firstMediaStreamTimestamp = performance.now() / 1000;
this._timestampOffset = -firstVideoFrameTimestamp;
} else {
this._timestampOffset = (performance.now() / 1000 - muxer.firstMediaStreamTimestamp)
- firstVideoFrameTimestamp;
}
}
if (this._encoder.getQueueSize() >= 4) {
// Drop frames if the encoder is overloaded
videoFrame.close();
return;
}
if (this._encoder.getQueueSize() >= 4) {
// Drop frames if the encoder is overloaded
videoFrame.close();
return;
}
void this._encoder.add(new VideoSample(videoFrame), true)
.catch((error) => {
this._abortController?.abort();
throw error;
});
},
});
void this._encoder.add(new VideoSample(videoFrame), true)
.catch((error) => {
errored = true;
processor.readable.pipeTo(consumer, {
signal: this._abortController.signal,
}).catch((err) => {
// Handle abort error silently
if (err instanceof DOMException && err.name === 'AbortError') return;
// Handle other errors
console.error('Pipe error:', err);
});
this._abortController?.abort();
this._promiseWithResolvers.reject(error);
if (this._workerTrackId !== null) {
// Tell the worker to stop the track
sendMessageToMediaStreamTrackProcessorWorker({
type: 'stopTrack',
trackId: this._workerTrackId,
});
}
});
};
if (typeof MediaStreamTrackProcessor !== 'undefined') {
// We can do it here directly, perfect
const processor = new MediaStreamTrackProcessor({ track: this._track });
const consumer = new WritableStream<VideoFrame>({ write: onVideoFrame });
processor.readable.pipeTo(consumer, {
signal: this._abortController.signal,
}).catch((error) => {
// Handle AbortError silently
if (error instanceof DOMException && error.name === 'AbortError') return;
this._promiseWithResolvers.reject(error);
});
} else {
// It might still be supported in a worker, so let's check that
const supportedInWorker = await mediaStreamTrackProcessorIsSupportedInWorker();
if (supportedInWorker) {
this._workerTrackId = nextMediaStreamTrackProcessorWorkerId++;
sendMessageToMediaStreamTrackProcessorWorker({
type: 'videoTrack',
trackId: this._workerTrackId,
track: this._track,
}, [this._track]);
this._workerListener = (event: MessageEvent) => {
const message = event.data as MediaStreamTrackProcessorWorkerMessage;
if (message.type === 'videoFrame' && message.trackId === this._workerTrackId) {
onVideoFrame(message.videoFrame);
} else if (message.type === 'error' && message.trackId === this._workerTrackId) {
this._promiseWithResolvers.reject(message.error);
}
};
mediaStreamTrackProcessorWorker!.addEventListener('message', this._workerListener);
} else {
throw new Error('MediaStreamTrackProcessor is required but not supported by this browser.');
}
}
}
/** @internal */
@@ -637,6 +720,32 @@ export class MediaStreamVideoTrackSource extends VideoSource {
this._abortController = null;
}
if (this._workerTrackId !== null) {
assert(this._workerListener);
sendMessageToMediaStreamTrackProcessorWorker({
type: 'stopTrack',
trackId: this._workerTrackId,
});
// Wait for the worker to stop the track
await new Promise<void>((resolve) => {
const listener = (event: MessageEvent) => {
const message = event.data as MediaStreamTrackProcessorWorkerMessage;
if (message.type === 'trackStopped' && message.trackId === this._workerTrackId) {
assert(this._workerListener);
mediaStreamTrackProcessorWorker!.removeEventListener('message', this._workerListener);
mediaStreamTrackProcessorWorker!.removeEventListener('message', listener);
resolve();
}
};
mediaStreamTrackProcessorWorker!.addEventListener('message', listener);
});
}
await this._encoder.flushAndClose();
}
}
@@ -783,78 +892,86 @@ class AudioEncoderWrapper {
constructor(private source: AudioSource, private encodingConfig: AudioEncodingConfig) {}
async add(audioSample: AudioSample, shouldClose: boolean) {
this.checkForEncoderError();
this.source._ensureValidAdd();
try {
this.checkForEncoderError();
this.source._ensureValidAdd();
// Ensure audio parameters remain constant
if (this.lastNumberOfChannels !== null && this.lastSampleRate !== null) {
if (
audioSample.numberOfChannels !== this.lastNumberOfChannels
|| audioSample.sampleRate !== this.lastSampleRate
) {
throw new Error(
`Audio parameters must remain constant. Expected ${this.lastNumberOfChannels} channels at`
+ ` ${this.lastSampleRate} Hz, got ${audioSample.numberOfChannels} channels at`
+ ` ${audioSample.sampleRate} Hz.`,
);
}
} else {
this.lastNumberOfChannels = audioSample.numberOfChannels;
this.lastSampleRate = audioSample.sampleRate;
}
if (!this.encoderInitialized) {
if (!this.ensureEncoderPromise) {
void this.ensureEncoder(audioSample);
// Ensure audio parameters remain constant
if (this.lastNumberOfChannels !== null && this.lastSampleRate !== null) {
if (
audioSample.numberOfChannels !== this.lastNumberOfChannels
|| audioSample.sampleRate !== this.lastSampleRate
) {
throw new Error(
`Audio parameters must remain constant. Expected ${this.lastNumberOfChannels} channels at`
+ ` ${this.lastSampleRate} Hz, got ${audioSample.numberOfChannels} channels at`
+ ` ${audioSample.sampleRate} Hz.`,
);
}
} else {
this.lastNumberOfChannels = audioSample.numberOfChannels;
this.lastSampleRate = audioSample.sampleRate;
}
// No, this "if" statement is not useless. Sometimes, the above call to `ensureEncoder` might have
// synchronously completed and the encoder is already initialized. In this case, we don't need to await the
// promise anymore. This also fixes nasty async race condition bugs when multiple code paths are calling
// this method: It's important that the call that initialized the encoder go through this code first.
if (!this.encoderInitialized) {
await this.ensureEncoderPromise;
if (!this.ensureEncoderPromise) {
void this.ensureEncoder(audioSample);
}
// No, this "if" statement is not useless. Sometimes, the above call to `ensureEncoder` might have
// synchronously completed and the encoder is already initialized. In this case, we don't need to await
// the promise anymore. This also fixes nasty async race condition bugs when multiple code paths are
// calling this method: It's important that the call that initialized the encoder go through this
// code first.
if (!this.encoderInitialized) {
await this.ensureEncoderPromise;
}
}
}
assert(this.encoderInitialized);
assert(this.encoderInitialized);
if (this.customEncoder) {
this.customEncoderQueueSize++;
const promise = this.customEncoderCallSerializer
.call(() => this.customEncoder!.encode(audioSample))
.then(() => {
this.customEncoderQueueSize--;
if (this.customEncoder) {
this.customEncoderQueueSize++;
const promise = this.customEncoderCallSerializer
.call(() => this.customEncoder!.encode(audioSample))
.then(() => {
this.customEncoderQueueSize--;
if (shouldClose) {
audioSample.close();
}
})
.catch((error: Error) => {
this.encoderError ??= error;
});
if (shouldClose) {
audioSample.close();
}
})
.catch((error: Error) => {
this.encoderError ??= error;
});
if (this.customEncoderQueueSize >= 4) {
await promise;
if (this.customEncoderQueueSize >= 4) {
await promise;
}
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
} else if (this.isPcmEncoder) {
await this.doPcmEncoding(audioSample, shouldClose);
} else {
assert(this.encoder);
const audioData = audioSample.toAudioData();
this.encoder.encode(audioData);
audioData.close();
if (shouldClose) {
audioSample.close();
}
if (this.encoder.encodeQueueSize >= 4) {
await new Promise(resolve => this.encoder!.addEventListener('dequeue', resolve, { once: true }));
}
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
}
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
} else if (this.isPcmEncoder) {
await this.doPcmEncoding(audioSample, shouldClose);
} else {
assert(this.encoder);
const audioData = audioSample.toAudioData();
this.encoder.encode(audioData);
audioData.close();
} finally {
if (shouldClose) {
// Make sure it's always closed, even if there was an error
audioSample.close();
}
if (this.encoder.encodeQueueSize >= 4) {
await new Promise(resolve => this.encoder!.addEventListener('dequeue', resolve, { once: true }));
}
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
}
}
@@ -1186,7 +1303,7 @@ export class AudioBufferSource extends AudioSource {
/** @internal */
private _encoder: AudioEncoderWrapper;
/** @internal */
private _accumulatedFrameCount = 0;
private _accumulatedTime = 0;
constructor(encodingConfig: AudioEncodingConfig) {
validateAudioEncodingConfig(encodingConfig);
@@ -1208,47 +1325,10 @@ export class AudioBufferSource extends AudioSource {
throw new TypeError('audioBuffer must be an AudioBuffer.');
}
const MAX_FLOAT_COUNT = 64 * 1024 * 1024;
const audioSamples = AudioSample.fromAudioBuffer(audioBuffer, this._accumulatedTime);
const promises = audioSamples.map(sample => this._encoder.add(sample, true));
const numberOfChannels = audioBuffer.numberOfChannels;
const sampleRate = audioBuffer.sampleRate;
const totalFrames = audioBuffer.length;
const maxFramesPerChunk = Math.floor(MAX_FLOAT_COUNT / numberOfChannels);
let currentRelativeFrame = 0;
let remainingFrames = totalFrames;
const promises: Promise<void>[] = [];
// Create AudioData in a chunked fashion so we don't create huge Float32Arrays
while (remainingFrames > 0) {
const framesToCopy = Math.min(maxFramesPerChunk, remainingFrames);
const chunkData = new Float32Array(numberOfChannels * framesToCopy);
for (let channel = 0; channel < numberOfChannels; channel++) {
audioBuffer.copyFromChannel(
chunkData.subarray(channel * framesToCopy, channel * framesToCopy + framesToCopy),
channel,
currentRelativeFrame,
);
}
const audioSample = new AudioSample({
format: 'f32-planar',
sampleRate,
numberOfFrames: framesToCopy,
numberOfChannels,
timestamp: (this._accumulatedFrameCount + currentRelativeFrame) / sampleRate,
data: chunkData,
});
promises.push(this._encoder.add(audioSample, true));
currentRelativeFrame += framesToCopy;
remainingFrames -= framesToCopy;
}
this._accumulatedFrameCount += totalFrames;
this._accumulatedTime += audioBuffer.duration;
return Promise.all(promises);
}
@@ -1272,6 +1352,20 @@ export class MediaStreamAudioTrackSource extends AudioSource {
private _abortController: AbortController | null = null;
/** @internal */
private _track: MediaStreamAudioTrack;
/** @internal */
private _audioContext: AudioContext | null = null;
/** @internal */
private _scriptProcessorNode: ScriptProcessorNode | null = null; // Deprecated but goated
/** @internal */
private _promiseWithResolvers = promiseWithResolvers();
/** @internal */
private _errorPromiseAccessed = false;
/** A promise that rejects upon any error within this source. This promise never resolves. */
get errorPromise() {
this._errorPromiseAccessed = true;
return this._promiseWithResolvers.promise;
}
constructor(track: MediaStreamAudioTrack, encodingConfig: AudioEncodingConfig) {
if (!(track instanceof MediaStreamTrack) || track.kind !== 'audio') {
@@ -1285,41 +1379,104 @@ export class MediaStreamAudioTrackSource extends AudioSource {
}
/** @internal */
override _start() {
override async _start() {
if (!this._errorPromiseAccessed) {
console.warn(
'Make sure not to ignore the `errorPromise` field on MediaStreamVideoTrackSource, so that any internal'
+ ' errors get bubbled up properly.',
);
}
this._abortController = new AbortController();
let dataReceived = false;
if (typeof MediaStreamTrackProcessor !== 'undefined') {
// Great, MediaStreamTrackProcessor is supported, this is the preferred way of doing things
let firstAudioDataTimestamp: number | null = null;
const processor = new MediaStreamTrackProcessor({ track: this._track });
const consumer = new WritableStream<AudioData>({
write: (audioData) => {
if (!dataReceived) {
setMediaStreamTimestampOffset(this, audioData);
dataReceived = true;
const processor = new MediaStreamTrackProcessor({ track: this._track });
const consumer = new WritableStream<AudioData>({
write: (audioData) => {
if (firstAudioDataTimestamp === null) {
firstAudioDataTimestamp = audioData.timestamp / 1e6;
const muxer = this._connectedTrack!.output._muxer;
if (muxer.firstMediaStreamTimestamp === null) {
muxer.firstMediaStreamTimestamp = performance.now() / 1000;
this._timestampOffset = -firstAudioDataTimestamp;
} else {
this._timestampOffset = (performance.now() / 1000 - muxer.firstMediaStreamTimestamp)
- firstAudioDataTimestamp;
}
}
if (this._encoder.getQueueSize() >= 4) {
// Drop data if the encoder is overloaded
audioData.close();
return;
}
void this._encoder.add(new AudioSample(audioData), true)
.catch((error) => {
this._abortController?.abort();
this._promiseWithResolvers.reject(error);
});
},
});
processor.readable.pipeTo(consumer, {
signal: this._abortController.signal,
}).catch((error) => {
// Handle AbortError silently
if (error instanceof DOMException && error.name === 'AbortError') return;
this._promiseWithResolvers.reject(error);
});
} else {
// Let's fall back to an AudioContext approach
this._audioContext = new AudioContext({ sampleRate: this._track.getSettings().sampleRate });
const sourceNode = this._audioContext.createMediaStreamSource(new MediaStream([this._track]));
this._scriptProcessorNode = this._audioContext.createScriptProcessor(4096);
if (this._audioContext.state === 'suspended') {
await this._audioContext.resume();
}
sourceNode.connect(this._scriptProcessorNode);
this._scriptProcessorNode.connect(this._audioContext.destination);
let audioReceived = false;
let totalDuration = 0;
this._scriptProcessorNode.onaudioprocess = (event) => {
const audioSamples = AudioSample.fromAudioBuffer(event.inputBuffer, totalDuration);
totalDuration += event.inputBuffer.duration;
for (const audioSample of audioSamples) {
if (!audioReceived) {
audioReceived = true;
const muxer = this._connectedTrack!.output._muxer;
if (muxer.firstMediaStreamTimestamp === null) {
muxer.firstMediaStreamTimestamp = performance.now() / 1000;
} else {
this._timestampOffset = performance.now() / 1000 - muxer.firstMediaStreamTimestamp;
}
}
if (this._encoder.getQueueSize() >= 4) {
// Drop data if the encoder is overloaded
audioSample.close();
continue;
}
void this._encoder.add(audioSample, true)
.catch((error) => {
void this._audioContext!.suspend();
this._promiseWithResolvers.reject(error);
});
}
if (this._encoder.getQueueSize() >= 4) {
// Drop data if the encoder is overloaded
audioData.close();
return;
}
void this._encoder.add(new AudioSample(audioData), true)
.catch((error) => {
this._abortController?.abort();
throw error;
});
},
});
processor.readable.pipeTo(consumer, {
signal: this._abortController.signal,
}).catch((err) => {
// Handle abort error silently
if (err instanceof DOMException && err.name === 'AbortError') return;
// Handle other errors
console.error('Pipe error:', err);
});
};
}
}
/** @internal */
@@ -1329,22 +1486,175 @@ export class MediaStreamAudioTrackSource extends AudioSource {
this._abortController = null;
}
if (this._audioContext) {
assert(this._scriptProcessorNode);
this._scriptProcessorNode.disconnect();
await this._audioContext.suspend();
}
await this._encoder.flushAndClose();
}
}
const setMediaStreamTimestampOffset = (source: MediaSource, sample: VideoFrame | AudioData) => {
const timestampInSeconds = sample.timestamp / 1e6;
// === MEDIA STREAM TRACK PROCESSOR WORKER ===
assert(source._connectedTrack);
const muxer = source._connectedTrack.output._muxer;
if (muxer.firstMediaStreamTimestamp === null) {
// We're the first MediaStreamTrack of this output to receive data
muxer.firstMediaStreamTimestamp = timestampInSeconds;
type MediaStreamTrackProcessorWorkerMessage = {
type: 'support';
supported: boolean;
} | {
type: 'videoFrame';
trackId: number;
videoFrame: VideoFrame;
} | {
type: 'trackStopped';
trackId: number;
} | {
type: 'error';
trackId: number;
error: Error;
};
type MediaStreamTrackProcessorControllerMessage = {
type: 'videoTrack';
trackId: number;
track: MediaStreamVideoTrack;
} | {
type: 'stopTrack';
trackId: number;
};
const mediaStreamTrackProcessorWorkerCode = () => {
const sendMessage = (message: MediaStreamTrackProcessorWorkerMessage, transfer?: Transferable[]) => {
if (transfer) {
// The error is bullshit, it's using the wrong postMessage
// eslint-disable-next-line @typescript-eslint/no-explicit-any, @typescript-eslint/no-unsafe-argument
self.postMessage(message, transfer as any);
} else {
self.postMessage(message);
}
};
// Immediately send a message to the main thread, letting them know of the support
sendMessage({
type: 'support',
supported: typeof MediaStreamTrackProcessor !== 'undefined',
});
const abortControllers = new Map<number, AbortController>();
const stoppedTracks = new Set<number>();
self.addEventListener('message', (event) => {
const message = event.data as MediaStreamTrackProcessorControllerMessage;
switch (message.type) {
case 'videoTrack': {
const processor = new MediaStreamTrackProcessor({ track: message.track });
const consumer = new WritableStream<VideoFrame>({
write: (videoFrame) => {
if (stoppedTracks.has(message.trackId)) {
videoFrame.close();
return;
}
// Send it to the main thread
sendMessage({
type: 'videoFrame',
trackId: message.trackId,
videoFrame,
}, [videoFrame]);
},
});
const abortController = new AbortController();
abortControllers.set(message.trackId, abortController);
processor.readable.pipeTo(consumer, {
signal: abortController.signal,
}).catch((error: Error) => {
// Handle AbortError silently
if (error instanceof DOMException && error.name === 'AbortError') return;
sendMessage({
type: 'error',
trackId: message.trackId,
error,
});
});
}; break;
case 'stopTrack': {
const abortController = abortControllers.get(message.trackId);
if (abortController) {
abortController.abort();
abortControllers.delete(message.trackId);
}
stoppedTracks.add(message.trackId);
sendMessage({
type: 'trackStopped',
trackId: message.trackId,
});
}; break;
default: assertNever(message);
}
});
};
let nextMediaStreamTrackProcessorWorkerId = 0;
let mediaStreamTrackProcessorWorker: Worker | null = null;
const initMediaStreamTrackProcessorWorker = () => {
const blob = new Blob(
[`(${mediaStreamTrackProcessorWorkerCode.toString()})()`],
{ type: 'application/javascript' },
);
const url = URL.createObjectURL(blob);
mediaStreamTrackProcessorWorker = new Worker(url);
};
let mediaStreamTrackProcessorIsSupportedInWorkerCache: boolean | null = null;
const mediaStreamTrackProcessorIsSupportedInWorker = async () => {
if (mediaStreamTrackProcessorIsSupportedInWorkerCache !== null) {
return mediaStreamTrackProcessorIsSupportedInWorkerCache;
}
// Math.min to ensure the timestamps can't get negative
source._timestampOffset = -Math.min(muxer.firstMediaStreamTimestamp, timestampInSeconds);
if (!mediaStreamTrackProcessorWorker) {
initMediaStreamTrackProcessorWorker();
}
return new Promise<boolean>((resolve) => {
assert(mediaStreamTrackProcessorWorker);
const listener = (event: MessageEvent) => {
const message = event.data as MediaStreamTrackProcessorWorkerMessage;
if (message.type === 'support') {
mediaStreamTrackProcessorIsSupportedInWorkerCache = message.supported;
mediaStreamTrackProcessorWorker!.removeEventListener('message', listener);
resolve(message.supported);
}
};
mediaStreamTrackProcessorWorker.addEventListener('message', listener);
});
};
const sendMessageToMediaStreamTrackProcessorWorker = (
message: MediaStreamTrackProcessorControllerMessage,
transfer?: Transferable[],
) => {
assert(mediaStreamTrackProcessorWorker);
if (transfer) {
mediaStreamTrackProcessorWorker.postMessage(message, transfer);
} else {
mediaStreamTrackProcessorWorker.postMessage(message);
}
};
/**