mirror of
https://github.com/arcodange-org/mediabunny.git
synced 2026-10-04 06:13:46 +02:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac8baa4873 | ||
|
|
232d1a6cd7 | ||
|
|
50fe065852 | ||
|
|
d38ad22559 | ||
|
|
9cc38329f2 | ||
|
|
09ed583b78 | ||
|
|
9de5b24ec5 |
@@ -22,9 +22,13 @@ Mediabunny is a JavaScript library for reading, writing, and converting media fi
|
||||
<a href="https://diffusion.studio/" target="_blank">
|
||||
<img src="./docs/public/sponsors/diffusionstudio.png" width="60" height="60" alt="Diffusion Studio">
|
||||
</a>
|
||||
|
||||
<a href="https://kino.ai/" target="_blank">
|
||||
<img src="./docs/public/sponsors/kino.jpg" width="60" height="60" alt="Kino">
|
||||
</a>
|
||||
</div>
|
||||
|
||||
[Get featured](https://github.com/sponsors/Vanilagy)
|
||||
[Sponsor Mediabunny's development](https://github.com/sponsors/Vanilagy)
|
||||
|
||||
## Features
|
||||
|
||||
|
||||
+13
-5
@@ -23,21 +23,29 @@
|
||||
format: new Mediabunny.Mp4OutputFormat(),
|
||||
});
|
||||
if (videoTrack) {
|
||||
output.addVideoTrack(new Mediabunny.MediaStreamVideoTrackSource(videoTrack, {
|
||||
const source = new Mediabunny.MediaStreamVideoTrackSource(videoTrack, {
|
||||
codec: 'avc',
|
||||
bitrate: Mediabunny.QUALITY_MEDIUM
|
||||
}));
|
||||
});
|
||||
|
||||
source.errorPromise.catch((d) => console.log("Hello?????", d));
|
||||
|
||||
output.addVideoTrack(source);
|
||||
}
|
||||
if (audioTrack) {
|
||||
output.addAudioTrack(new Mediabunny.MediaStreamAudioTrackSource(audioTrack, {
|
||||
const source = new Mediabunny.MediaStreamAudioTrackSource(audioTrack, {
|
||||
codec: 'aac',
|
||||
bitrate: Mediabunny.QUALITY_MEDIUM
|
||||
}));
|
||||
});
|
||||
|
||||
source.errorPromise.catch((d) => console.log("Hello!!???", d));
|
||||
|
||||
output.addAudioTrack(source);
|
||||
}
|
||||
|
||||
await output.start();
|
||||
|
||||
await new Promise(resolve => setTimeout(resolve, 3000));
|
||||
await new Promise(resolve => setTimeout(resolve, 5000));
|
||||
|
||||
await output.finalize();
|
||||
|
||||
|
||||
@@ -163,6 +163,9 @@ const videoTrackSource = new MediaStreamVideoTrackSource(videoTrack, {
|
||||
codec: 'vp9',
|
||||
bitrate: 1e7,
|
||||
});
|
||||
|
||||
// Make sure to allow any internal errors to properly bubble up
|
||||
videoTrackSource.errorPromise.catch((error) => ...);
|
||||
```
|
||||
|
||||
This source requires no additional method calls; data will automatically be captured and piped to the output file as soon as `start()` is called on the `Output`. Make sure to `stop()` on `videoTrack` after finalizing the `Output` if you don't need the user's media anymore.
|
||||
@@ -171,6 +174,10 @@ This source requires no additional method calls; data will automatically be capt
|
||||
If this source is the only MediaStreamTrack source in the `Output`, then the first video sample added by it starts at timestamp 0. If there are multiple, then the earliest media sample across all tracks starts at timestamp 0, and all tracks will be perfectly synchronized with each other.
|
||||
:::
|
||||
|
||||
::: warning
|
||||
`MediaStreamVideoTrackSource`'s internals are detached from the typical code flow but can still throw, so make sure to utilize `errorPromise` to deal with any errors and to stop the `Output`.
|
||||
:::
|
||||
|
||||
### `EncodedVideoPacketSource`
|
||||
|
||||
The most barebones of all video sources, this source can be used to directly pipe [encoded packets](./packets-and-samples#encodedpacket) of video data to the output. This source requires that you take care of the encoding process yourself, which enables you to use the WebCodecs API manually or to plug in your own encoding stack. Alternatively, you may retrieve the encoded packets directly by reading them from another media file, allowing you to skip decoding and reencoding video data.
|
||||
@@ -312,6 +319,9 @@ const audioTrackSource = new MediaStreamAudioTrackSource(audioTrack, {
|
||||
codec: 'opus',
|
||||
bitrate: 128e3,
|
||||
});
|
||||
|
||||
// Make sure to allow any internal errors to properly bubble up
|
||||
audioTrackSource.errorPromise.catch((error) => ...);
|
||||
```
|
||||
|
||||
This source requires no additional method calls; data will automatically be captured and piped to the output file as soon as `start()` is called on the `Output`. Make sure to `stop()` on `audioTrack` after finalizing the `Output` if you don't need the user's media anymore.
|
||||
@@ -320,6 +330,10 @@ This source requires no additional method calls; data will automatically be capt
|
||||
If this source is the only MediaStreamTrack source in the `Output`, then the first audio sample added by it starts at timestamp 0. If there are multiple, then the earliest media sample across all tracks starts at timestamp 0, and all tracks will be perfectly synchronized with each other.
|
||||
:::
|
||||
|
||||
::: warning
|
||||
`MediaStreamAudioTrackSource`'s internals are detached from the typical code flow but can still throw, so make sure to utilize `errorPromise` to deal with any errors and to stop the `Output`.
|
||||
:::
|
||||
|
||||
### `EncodedAudioPacketSource`
|
||||
|
||||
The most barebones of all audio sources, this source can be used to directly pipe [encoded packets](./packets-and-samples#encodedpacket) of audio data to the output. This source requires that you take care of the encoding process yourself, which enables you to use the WebCodecs API manually or to plug in your own encoding stack. Alternatively, you may retrieve the encoded packets directly by reading them from another media file, allowing you to skip decoding and reencoding audio data.
|
||||
|
||||
@@ -401,7 +401,7 @@ An audio sample represents a section of audio data. It can be created directly f
|
||||
|
||||
### Creating audio samples
|
||||
|
||||
Audio samples can be constructed either from an `AudioData` instance or an initialization object:
|
||||
Audio samples can be constructed either from an `AudioData` instance, an initialization object, or an `AudioBuffer`:
|
||||
|
||||
```ts
|
||||
import { AudioSample } from 'mediabunny';
|
||||
@@ -417,6 +417,11 @@ const sample = new AudioSample({
|
||||
sampleRate: 44100, // in Hz
|
||||
timestamp: 0, // in seconds
|
||||
});
|
||||
|
||||
// From AudioBuffer:
|
||||
const timestamp = 0; // in seconds
|
||||
const samples = AudioSample.fromAudioBuffer(audioBuffer, timestamp);
|
||||
// => Returns multiple AudioSamples if the AudioBuffer is very long
|
||||
```
|
||||
|
||||
The following audio sample formats are supported:
|
||||
|
||||
@@ -90,9 +90,11 @@ const sponsors = {
|
||||
gold: [
|
||||
{ image: '/sponsors/gling.svg', name: 'Gling AI', url: 'https://www.gling.ai/' },
|
||||
{ image: '/sponsors/diffusionstudio.png', name: 'Diffusion Studio', url: 'https://diffusion.studio/' },
|
||||
{ image: '/sponsors/kino.jpg', name: 'Kino', url: 'https://kino.ai/' },
|
||||
],
|
||||
individual: [
|
||||
{ image: 'https://avatars.githubusercontent.com/u/84167135', name: 'Memenome', url: 'https://github.com/memenome' },
|
||||
{ image: 'https://avatars.githubusercontent.com/u/9549394', name: 'studnitz', url: 'https://github.com/studnitz' },
|
||||
{ image: 'https://avatars.githubusercontent.com/u/30229596', name: 'Pablo Bonilla', url: 'https://github.com/devPablo' },
|
||||
{ image: 'https://avatars.githubusercontent.com/u/58149663', name: 'H7GhosT', url: 'https://github.com/H7GhosT' },
|
||||
{ image: 'https://avatars.githubusercontent.com/u/91711202', name: 'ihasq', url: 'https://github.com/ihasq' },
|
||||
|
||||
Binary file not shown.
|
After Width: | Height: | Size: 6.0 KiB |
@@ -98,6 +98,10 @@ const compressFile = async (file: File) => {
|
||||
compressionFacts.textContent
|
||||
= `${(output.target.buffer!.byteLength / file.size * 100).toPrecision(3)}% of original size`;
|
||||
} catch (error) {
|
||||
console.error(error);
|
||||
|
||||
await currentConversion?.cancel();
|
||||
|
||||
errorElement.textContent = String(error);
|
||||
clearInterval(currentIntervalId);
|
||||
|
||||
|
||||
@@ -22,6 +22,7 @@
|
||||
<hr class="w-full max-w-96 my-4 border-zinc-300 dark:border-zinc-700" style="display: none;">
|
||||
|
||||
<p id="error-element" class="text-red-500"></p>
|
||||
<p id="warning-element" class="text-amber-500"></p>
|
||||
|
||||
<div class="flex gap-4" id="main-container" style="display: none;">
|
||||
<div class="flex flex-col items-center">
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import {
|
||||
canEncodeAudio,
|
||||
CanvasSource,
|
||||
MediaStreamAudioTrackSource,
|
||||
Mp4OutputFormat,
|
||||
@@ -13,6 +14,7 @@ const mainContainer = document.querySelector('#main-container') as HTMLDivElemen
|
||||
const videoElement = document.querySelector('video') as HTMLVideoElement;
|
||||
const downloadButton = document.querySelector('#download-button') as HTMLAnchorElement;
|
||||
const errorElement = document.querySelector('#error-element') as HTMLParagraphElement;
|
||||
const warningElement = document.querySelector('#warning-element') as HTMLParagraphElement;
|
||||
|
||||
const canvas = document.querySelector('canvas') as HTMLCanvasElement;
|
||||
const context = canvas.getContext('2d', { alpha: false, desynchronized: true })!;
|
||||
@@ -39,19 +41,29 @@ const startRecording = async () => {
|
||||
videoElement.src = '';
|
||||
downloadButton.style.display = 'none';
|
||||
errorElement.textContent = '';
|
||||
warningElement.textContent = '';
|
||||
|
||||
// Paint a white background to the canvas
|
||||
context.fillStyle = 'white';
|
||||
context.fillRect(0, 0, canvas.width, canvas.height);
|
||||
|
||||
// Get user microphone
|
||||
mediaStream = await navigator.mediaDevices.getUserMedia({ audio: true });
|
||||
const audioIsEncodable = await canEncodeAudio('opus', {
|
||||
bitrate: QUALITY_MEDIUM,
|
||||
});
|
||||
|
||||
let audioTrack: MediaStreamAudioTrack | null = null;
|
||||
if (audioIsEncodable) {
|
||||
// Get user microphone
|
||||
mediaStream = await navigator.mediaDevices.getUserMedia({ audio: true });
|
||||
audioTrack = mediaStream.getAudioTracks()[0] ?? null;
|
||||
} else {
|
||||
warningElement.textContent
|
||||
= 'Audio is not yet encodable by your browser, so the audio track has been omitted.';
|
||||
}
|
||||
|
||||
horizontalRule.style.display = '';
|
||||
mainContainer.style.display = '';
|
||||
|
||||
const audioTrack = mediaStream.getAudioTracks()[0];
|
||||
|
||||
// Create a new output file
|
||||
output = new Output({
|
||||
// We're using fragmented MP4 here; streamable WebM would also work
|
||||
@@ -99,6 +111,8 @@ const startRecording = async () => {
|
||||
codec: 'opus',
|
||||
bitrate: QUALITY_MEDIUM,
|
||||
});
|
||||
audioSource.errorPromise.catch(cancelRecording); // Make sure errors are bubbled up
|
||||
|
||||
output.addAudioTrack(audioSource);
|
||||
}
|
||||
|
||||
@@ -108,9 +122,9 @@ const startRecording = async () => {
|
||||
readyForMoreFrames = true;
|
||||
lastFrameNumber = -1;
|
||||
|
||||
// Start the video frame capture loop
|
||||
void addVideoFrame();
|
||||
videoCaptureInterval = window.setInterval(() => void addVideoFrame(), 1000 / frameRate);
|
||||
// Start the video frame capture loop, making sure errors are caught
|
||||
void addVideoFrame().catch(cancelRecording);
|
||||
videoCaptureInterval = window.setInterval(() => void addVideoFrame().catch(cancelRecording), 1000 / frameRate);
|
||||
|
||||
const mimeType = await output.getMimeType();
|
||||
sourceBuffer = mediaSource.addSourceBuffer(mimeType);
|
||||
@@ -121,21 +135,35 @@ const startRecording = async () => {
|
||||
toggleRecordingButton.textContent = 'Stop recording';
|
||||
toggleRecordingButton.disabled = false;
|
||||
} catch (error) {
|
||||
errorElement.textContent = String(error);
|
||||
|
||||
mainContainer.style.display = 'none';
|
||||
toggleRecordingButton.textContent = 'Start recording';
|
||||
toggleRecordingButton.disabled = false;
|
||||
recording = false;
|
||||
await cancelRecording(error);
|
||||
}
|
||||
};
|
||||
|
||||
const cancelRecording = async (error: unknown) => {
|
||||
if (!recording) {
|
||||
return; // Already canceled
|
||||
}
|
||||
|
||||
console.error(error);
|
||||
|
||||
errorElement.textContent = String(error);
|
||||
|
||||
clearInterval(videoCaptureInterval);
|
||||
mainContainer.style.display = 'none';
|
||||
toggleRecordingButton.textContent = 'Start recording';
|
||||
toggleRecordingButton.disabled = false;
|
||||
recording = false;
|
||||
await output?.cancel();
|
||||
|
||||
mediaStream?.getTracks().forEach(track => track.stop());
|
||||
};
|
||||
|
||||
const stopRecording = async () => {
|
||||
toggleRecordingButton.textContent = 'Stopping...';
|
||||
toggleRecordingButton.disabled = true;
|
||||
|
||||
clearInterval(videoCaptureInterval);
|
||||
mediaStream.getTracks().forEach(track => track.stop());
|
||||
mediaStream?.getTracks().forEach(track => track.stop());
|
||||
|
||||
await output.finalize();
|
||||
|
||||
|
||||
@@ -175,8 +175,10 @@ const initMediaPlayer = async (file: File) => {
|
||||
controlsElement.style.opacity = '1';
|
||||
playerContainer.style.cursor = '';
|
||||
}
|
||||
} catch (e) {
|
||||
errorElement.textContent = String(e);
|
||||
} catch (error) {
|
||||
console.error(error);
|
||||
|
||||
errorElement.textContent = String(error);
|
||||
playerContainer.style.display = 'none';
|
||||
}
|
||||
};
|
||||
|
||||
@@ -129,6 +129,8 @@ const renderObject = (object: Record<string, unknown>) => {
|
||||
listItem.removeChild(loadingSpan);
|
||||
listItem.appendChild(renderValue(resolvedValue));
|
||||
}).catch((error) => {
|
||||
console.error(error);
|
||||
|
||||
// Show the promise error
|
||||
listItem.removeChild(loadingSpan);
|
||||
const errorSpan = document.createElement('span');
|
||||
|
||||
@@ -7,6 +7,7 @@ import {
|
||||
QUALITY_HIGH,
|
||||
getFirstEncodableAudioCodec,
|
||||
getFirstEncodableVideoCodec,
|
||||
OutputFormat,
|
||||
} from 'mediabunny';
|
||||
|
||||
const durationSlider = document.querySelector('#duration-slider') as HTMLInputElement;
|
||||
@@ -57,6 +58,8 @@ let currentScaleIndex = 0;
|
||||
let collisionCount = 0;
|
||||
let collisionsPerScale = 0;
|
||||
|
||||
let output: Output<OutputFormat, BufferTarget>;
|
||||
|
||||
/** === MAIN VIDEO FILE GENERATION LOGIC === */
|
||||
|
||||
const generateVideo = async () => {
|
||||
@@ -82,7 +85,7 @@ const generateVideo = async () => {
|
||||
initScene(duration);
|
||||
|
||||
// Create a new output file
|
||||
const output = new Output({
|
||||
output = new Output({
|
||||
target: new BufferTarget(), // Stored in memory
|
||||
format: new Mp4OutputFormat(),
|
||||
});
|
||||
@@ -184,6 +187,10 @@ const generateVideo = async () => {
|
||||
const fileSizeMiB = (videoBlob.size / (1024 * 1024)).toPrecision(3);
|
||||
videoInfo.textContent = `File size: ${fileSizeMiB} MiB`;
|
||||
} catch (error) {
|
||||
console.error(error);
|
||||
|
||||
await output?.cancel();
|
||||
|
||||
clearInterval(progressInterval);
|
||||
errorElement.textContent = String(error);
|
||||
progressBarContainer.style.display = 'none';
|
||||
|
||||
@@ -96,8 +96,10 @@ const generateThumbnails = async (file: File) => {
|
||||
|
||||
i++;
|
||||
}
|
||||
} catch (e) {
|
||||
errorElement.textContent = String(e);
|
||||
} catch (error) {
|
||||
console.error(error);
|
||||
|
||||
errorElement.textContent = String(error);
|
||||
thumbnailContainer.innerHTML = '';
|
||||
}
|
||||
};
|
||||
|
||||
Generated
+2
-2
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "mediabunny",
|
||||
"version": "1.3.3",
|
||||
"version": "1.4.1",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "mediabunny",
|
||||
"version": "1.3.3",
|
||||
"version": "1.4.1",
|
||||
"license": "MPL-2.0",
|
||||
"dependencies": {
|
||||
"@types/dom-mediacapture-transform": "^0.1.11",
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "mediabunny",
|
||||
"author": "Vanilagy",
|
||||
"version": "1.3.3",
|
||||
"version": "1.4.1",
|
||||
"description": "Pure TypeScript media toolkit for reading, writing, and converting media files, directly in the browser.",
|
||||
"type": "module",
|
||||
"main": "./dist/bundles/mediabunny.cjs",
|
||||
|
||||
@@ -17,6 +17,8 @@ const checkDocblocks = (filePath: string) => {
|
||||
ts.isInterfaceDeclaration(node)
|
||||
|| ts.isClassDeclaration(node)
|
||||
|| ts.isMethodDeclaration(node)
|
||||
|| ts.isGetAccessorDeclaration(node)
|
||||
|| ts.isSetAccessorDeclaration(node)
|
||||
|| ts.isPropertyDeclaration(node)
|
||||
|| ts.isFunctionDeclaration(node)
|
||||
|| ts.isTypeAliasDeclaration(node)
|
||||
|
||||
@@ -32,6 +32,7 @@ import {
|
||||
transformAnnexBToLengthPrefixed,
|
||||
} from '../codec-data';
|
||||
import { buildIsobmffMimeType } from './isobmff-misc';
|
||||
import { MAX_BOX_HEADER_SIZE, MIN_BOX_HEADER_SIZE } from './isobmff-reader';
|
||||
|
||||
export const GLOBAL_TIMESCALE = 1000;
|
||||
const TIMESTAMP_OFFSET = 2_082_844_800; // Seconds between Jan 1 1904 and Jan 1 1970
|
||||
@@ -989,10 +990,7 @@ export class IsobmffMuxer extends Muxer {
|
||||
const moofOffset = this.writer.getPos();
|
||||
const mdatStartPos = moofOffset + this.boxWriter.measureBox(moofBox);
|
||||
|
||||
// Header with large size. We always reserve 16 bytes for it even if we don't end up using the large size.
|
||||
const mdatHeaderSize = 16;
|
||||
|
||||
let currentPos = mdatStartPos + mdatHeaderSize;
|
||||
let currentPos = mdatStartPos + MIN_BOX_HEADER_SIZE;
|
||||
let fragmentStartTimestamp = Infinity;
|
||||
for (const trackData of tracksInFragment) {
|
||||
trackData.currentChunk!.offset = currentPos;
|
||||
@@ -1006,6 +1004,15 @@ export class IsobmffMuxer extends Muxer {
|
||||
}
|
||||
|
||||
const mdatSize = currentPos - mdatStartPos;
|
||||
const needsLargeMdatSize = mdatSize >= 2 ** 32;
|
||||
|
||||
if (needsLargeMdatSize) {
|
||||
// Shift all offsets by 8. Previously, all chunks were shifted assuming the large box size, but due to what
|
||||
// I suspect is a bug in WebKit, it failed in Safari (when livestreaming with MSE, not for static playback).
|
||||
for (const trackData of tracksInFragment) {
|
||||
trackData.currentChunk!.offset! += MAX_BOX_HEADER_SIZE - MIN_BOX_HEADER_SIZE;
|
||||
}
|
||||
}
|
||||
|
||||
if (this.format._options.onMoof) {
|
||||
this.writer.startTrackingWrites();
|
||||
@@ -1025,11 +1032,11 @@ export class IsobmffMuxer extends Muxer {
|
||||
this.writer.startTrackingWrites();
|
||||
}
|
||||
|
||||
const mdatBox = mdat(mdatSize >= 2 ** 32);
|
||||
const mdatBox = mdat(needsLargeMdatSize);
|
||||
mdatBox.size = mdatSize;
|
||||
this.boxWriter.writeBox(mdatBox);
|
||||
|
||||
this.writer.seek(mdatStartPos + mdatHeaderSize);
|
||||
this.writer.seek(mdatStartPos + (needsLargeMdatSize ? MAX_BOX_HEADER_SIZE : MIN_BOX_HEADER_SIZE));
|
||||
|
||||
// Write sample data
|
||||
for (const trackData of tracksInFragment) {
|
||||
|
||||
@@ -368,6 +368,13 @@ export class MatroskaDemuxer extends Demuxer {
|
||||
this.readContiguousElements(this.metadataReader, size);
|
||||
}
|
||||
|
||||
if (this.currentSegment.timestampScale === -1) {
|
||||
// TimestampScale element is missing. Technically an invalid file, but let's default to the typical value,
|
||||
// which is 1e6.
|
||||
this.currentSegment.timestampScale = 1e6;
|
||||
this.currentSegment.timestampFactor = 1e9 / 1e6;
|
||||
}
|
||||
|
||||
// Put default tracks first
|
||||
this.currentSegment.tracks.sort((a, b) => Number(b.isDefault) - Number(a.isDefault));
|
||||
|
||||
|
||||
+545
-235
@@ -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);
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
|
||||
+2
-3
@@ -344,9 +344,8 @@ export class Output<
|
||||
|
||||
await this._muxer.start();
|
||||
|
||||
for (const track of this._tracks) {
|
||||
track.source._start();
|
||||
}
|
||||
const promises = this._tracks.map(track => track.source._start());
|
||||
await Promise.all(promises);
|
||||
|
||||
release();
|
||||
})();
|
||||
|
||||
@@ -1063,6 +1063,58 @@ export class AudioSample {
|
||||
// eslint-disable-next-line @typescript-eslint/no-unnecessary-type-assertion
|
||||
(this.timestamp as number) = newTimestamp;
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates AudioSamples from an AudioBuffer, starting at the given timestamp in seconds. Typically creates exactly
|
||||
* one sample, but may create multiple if the AudioBuffer is exceedingly large.
|
||||
*/
|
||||
static fromAudioBuffer(audioBuffer: AudioBuffer, timestamp: number) {
|
||||
if (!(audioBuffer instanceof AudioBuffer)) {
|
||||
throw new TypeError('audioBuffer must be an AudioBuffer.');
|
||||
}
|
||||
|
||||
const MAX_FLOAT_COUNT = 64 * 1024 * 1024;
|
||||
|
||||
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 result: AudioSample[] = [];
|
||||
|
||||
// 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: timestamp + currentRelativeFrame / sampleRate,
|
||||
data: chunkData,
|
||||
});
|
||||
|
||||
result.push(audioSample);
|
||||
|
||||
currentRelativeFrame += framesToCopy;
|
||||
remainingFrames -= framesToCopy;
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
||||
const getBytesPerSample = (format: AudioSampleFormat): number => {
|
||||
|
||||
Reference in New Issue
Block a user