mirror of
https://github.com/arcodange-org/mediabunny.git
synced 2026-10-10 01:03:45 +02:00
Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d12c48562c | ||
|
|
77e70f9090 | ||
|
|
9a7120f5f4 | ||
|
|
2ab6bbe3d1 | ||
|
|
670664c76a | ||
|
|
c8c8cf214c | ||
|
|
7d9c063105 |
@@ -115,7 +115,9 @@ await output.finalize();
|
||||
|
||||
### Upload to a server
|
||||
|
||||
This models a stream upload, where files are being uploaded *while* they are being created.
|
||||
#### Stream upload
|
||||
|
||||
This code models a stream upload, where files are being uploaded *while* they are being created:
|
||||
|
||||
```ts
|
||||
const promises: Promise<Response>[] = [];
|
||||
@@ -146,17 +148,52 @@ const output = new Output({
|
||||
return new StreamTarget(writable);
|
||||
},
|
||||
),
|
||||
onFinalize: () => Promise.all(promises),
|
||||
// ...
|
||||
});
|
||||
|
||||
// ...
|
||||
await output.finalize();
|
||||
|
||||
await Promise.all(promises);
|
||||
// All files have been uploaded to the server
|
||||
```
|
||||
|
||||
If this is too fancy, you can always use `BufferTarget` instead and upload its contents to a server in a non-streaming way after the `finalized` event.
|
||||
#### Monolithic upload
|
||||
|
||||
If streaming is not possible (e.g. when uploading to S3 via signed `PutObject`, which requires a known `Content-Length`), you can use [`BufferTarget`](../api/BufferTarget) with the [`onFinalize`](../api/BufferTargetOptions#onfinalize) option instead.
|
||||
|
||||
You could call `fetch` directly but this would halt Mediabunny's internals until the upload has completed. Instead, using a [`ConcurrentRunner`](../api/ConcurrentRunner) allows Mediabunny to keep producing data internally while the upload is in flight, while also allowing multiple concurrent uploads:
|
||||
|
||||
```ts
|
||||
import { ConcurrentRunner, ... } from 'mediabunny';
|
||||
|
||||
// This Mediabunny utility class is used to allow up to two requests
|
||||
// to run concurrently. When this number is exceeded, backpressure is
|
||||
// automatically applied internally.
|
||||
const runner = new ConcurrentRunner(2);
|
||||
|
||||
const output = new Output({
|
||||
target: new PathedTarget(
|
||||
'master.m3u8',
|
||||
({ path, mimeType }) =>
|
||||
new BufferTarget({
|
||||
onFinalize: buffer => runner.run(() =>
|
||||
fetch(`/upload?file=${encodeURIComponent(path)}`, {
|
||||
method: 'PUT',
|
||||
body: buffer,
|
||||
headers: {
|
||||
'Content-Type': mimeType,
|
||||
},
|
||||
})
|
||||
),
|
||||
}),
|
||||
),
|
||||
onFinalize: () => runner.flush(),
|
||||
// ...
|
||||
});
|
||||
|
||||
await output.finalize();
|
||||
// All files have been uploaded to the server
|
||||
```
|
||||
|
||||
## Adding tracks & media
|
||||
|
||||
@@ -224,6 +261,22 @@ You can extend this pattern to offer content in multiple codecs as well.
|
||||
|
||||
For the full list of transformation options, see [`VideoTransformOptions`](../api/VideoTransformOptions) and [`AudioTransformOptions`](../api/AudioTransformOptions).
|
||||
|
||||
---
|
||||
|
||||
If you're using the [Conversion API](./converting-media-files), you achieve the same thing using [output track fan-out](./converting-media-files#track-fan-out):
|
||||
```ts
|
||||
const conversion = await Conversion.init({
|
||||
input,
|
||||
output,
|
||||
video: [
|
||||
{ height: 1080, bitrate: QUALITY_VERY_HIGH },
|
||||
{ height: 720, bitrate: QUALITY_HIGH },
|
||||
{ height: 480, bitrate: QUALITY_MEDIUM },
|
||||
{ height: 360, bitrate: QUALITY_LOW },
|
||||
],
|
||||
});
|
||||
```
|
||||
|
||||
### Track metadata
|
||||
|
||||
Often you'll want to provide additional [track metadata](../api/BaseTrackMetadata) when dealing with multiple tracks. For example:
|
||||
|
||||
@@ -216,6 +216,10 @@ await output.finalize();
|
||||
const file = output.target.buffer; // => Uint8Array
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
An optional [`onFinalize`](../api/OutputOptions#onfinalize) callback can be provided in the output options. This function will be called at the end of `.finalize()` and can be used to do work once the output has been completed. If it returns a promise, it will be awaited by the output.
|
||||
|
||||
## Canceling an output
|
||||
|
||||
Sometimes, you may want to cancel the ongoing creation of an output file. For this, use the `cancel` method:
|
||||
@@ -288,6 +292,24 @@ output.target.buffer; // => ArrayBuffer
|
||||
|
||||
This target is a great choice for small-ish files (< 100 MB), but since all data will be kept in memory, using it for large files is suboptimal. If the output gets very large, the page might crash due to memory exhaustion. For these cases, using `StreamTarget` is recommended.
|
||||
|
||||
#### `onFinalize` callback
|
||||
|
||||
`BufferTarget` accepts an `onFinalize` option, which is called with the complete buffer once the target has been finalized. Useful for uploading the final buffer to a server or object store (e.g. S3 `PutObject`, which requires a known `Content-Length`):
|
||||
|
||||
```ts
|
||||
const output = new Output({
|
||||
target: new BufferTarget({
|
||||
onFinalize: async (buffer) => {
|
||||
await fetch('/upload', { method: 'PUT', body: buffer });
|
||||
},
|
||||
}),
|
||||
// ...
|
||||
});
|
||||
|
||||
await output.finalize();
|
||||
// The upload has completed by the time finalize resolves.
|
||||
```
|
||||
|
||||
### `StreamTarget`
|
||||
|
||||
This target passes you the data written by the `Output` in small chunks, requiring you to pipe that data elsewhere to manually assemble the final file. Example use cases include writing the file directly to disk, or uploading it to a server over the network.
|
||||
|
||||
Generated
+6
-6
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "mediabunny",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "mediabunny",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"license": "MPL-2.0",
|
||||
"workspaces": [
|
||||
"packages/*"
|
||||
@@ -12077,7 +12077,7 @@
|
||||
},
|
||||
"packages/aac-encoder": {
|
||||
"name": "@mediabunny/aac-encoder",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"license": "MPL-2.0",
|
||||
"devDependencies": {
|
||||
"@types/emscripten": "^1.40.1"
|
||||
@@ -12092,7 +12092,7 @@
|
||||
},
|
||||
"packages/ac3": {
|
||||
"name": "@mediabunny/ac3",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"license": "MPL-2.0",
|
||||
"devDependencies": {
|
||||
"@types/emscripten": "^1.40.1"
|
||||
@@ -12107,7 +12107,7 @@
|
||||
},
|
||||
"packages/flac-encoder": {
|
||||
"name": "@mediabunny/flac-encoder",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"license": "MPL-2.0",
|
||||
"devDependencies": {
|
||||
"@types/emscripten": "^1.40.1"
|
||||
@@ -12122,7 +12122,7 @@
|
||||
},
|
||||
"packages/mp3-encoder": {
|
||||
"name": "@mediabunny/mp3-encoder",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"license": "MPL-2.0",
|
||||
"devDependencies": {
|
||||
"@types/emscripten": "^1.40.1"
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "mediabunny",
|
||||
"author": "Vanilagy",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"description": "Pure TypeScript media toolkit for reading, writing, and converting media files, directly in the browser.",
|
||||
"type": "module",
|
||||
"workspaces": [
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "@mediabunny/aac-encoder",
|
||||
"author": "Vanilagy",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"description": "AAC encoder extension for Mediabunny, based on FFmpeg.",
|
||||
"main": "./dist/bundles/mediabunny-aac-encoder.mjs",
|
||||
"module": "./dist/bundles/mediabunny-aac-encoder.mjs",
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "@mediabunny/ac3",
|
||||
"author": "Vanilagy",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"description": "AC-3 and E-AC-3 (Dolby Digital) decoder and encoder extension for Mediabunny, based on FFmpeg.",
|
||||
"main": "./dist/bundles/mediabunny-ac3.mjs",
|
||||
"module": "./dist/bundles/mediabunny-ac3.mjs",
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "@mediabunny/flac-encoder",
|
||||
"author": "Vanilagy",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"description": "FLAC encoder extension for Mediabunny, based on libFLAC.",
|
||||
"main": "./dist/bundles/mediabunny-flac-encoder.mjs",
|
||||
"module": "./dist/bundles/mediabunny-flac-encoder.mjs",
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"name": "@mediabunny/mp3-encoder",
|
||||
"author": "Vanilagy",
|
||||
"version": "1.42.0-beta.3",
|
||||
"version": "1.42.0-beta.4",
|
||||
"description": "MP3 encoder extension for Mediabunny, based on LAME.",
|
||||
"main": "./dist/bundles/mediabunny-mp3-encoder.mjs",
|
||||
"module": "./dist/bundles/mediabunny-mp3-encoder.mjs",
|
||||
|
||||
@@ -16,7 +16,6 @@ const checkDocblocks = (filePath: string) => {
|
||||
if (
|
||||
ts.isInterfaceDeclaration(node)
|
||||
|| ts.isClassDeclaration(node)
|
||||
|| ts.isConstructorDeclaration(node)
|
||||
|| ts.isMethodDeclaration(node)
|
||||
|| ts.isGetAccessorDeclaration(node)
|
||||
|| ts.isSetAccessorDeclaration(node)
|
||||
@@ -61,13 +60,7 @@ const checkDocblocks = (filePath: string) => {
|
||||
let name = 'anonymous';
|
||||
const kind = ts.SyntaxKind[node.kind].replace(/Declaration|Statement/g, '').toLowerCase();
|
||||
|
||||
if (ts.isConstructorDeclaration(node)) {
|
||||
// For constructors, use the parent class name
|
||||
const parent = node.parent;
|
||||
if (ts.isClassDeclaration(parent) && parent.name) {
|
||||
name = parent.name.text;
|
||||
}
|
||||
} else if ('name' in node && node.name) {
|
||||
if ('name' in node && node.name) {
|
||||
if (ts.isIdentifier(node.name)) {
|
||||
name = node.name.text;
|
||||
} else if ('getText' in node.name) {
|
||||
|
||||
+67
-46
@@ -8,7 +8,16 @@
|
||||
|
||||
import { MediaCodec, validateAudioChunkMetadata, validateVideoChunkMetadata } from '../codec';
|
||||
import { EncodedAudioPacketSource, EncodedVideoPacketSource } from '../media-source';
|
||||
import { arrayArgmax, assert, findLastIndex, joinPaths, textEncoder, toArray, UNDETERMINED_LANGUAGE } from '../misc';
|
||||
import {
|
||||
arrayArgmax,
|
||||
assert,
|
||||
AsyncMutex,
|
||||
findLastIndex,
|
||||
joinPaths,
|
||||
textEncoder,
|
||||
toArray,
|
||||
UNDETERMINED_LANGUAGE,
|
||||
} from '../misc';
|
||||
import { Muxer } from '../muxer';
|
||||
import {
|
||||
Output,
|
||||
@@ -79,6 +88,11 @@ type Playlist = {
|
||||
nextOffset: number;
|
||||
info: HlsOutputSegmentInfo;
|
||||
} | null;
|
||||
|
||||
// For HLS, having a single mutex is too coarse. Every playlist is basically independent and therefore we can have
|
||||
// a per-playlist mutex instead of a per-muxer one. This means two packets from different playlists coming in don't
|
||||
// block each other.
|
||||
mutex: AsyncMutex;
|
||||
};
|
||||
|
||||
type PlaylistDeclaration = {
|
||||
@@ -131,6 +145,8 @@ export class HlsMuxer extends Muxer {
|
||||
}
|
||||
|
||||
async start(): Promise<void> {
|
||||
const release = await this.mutex.acquire();
|
||||
|
||||
const someRelative = this.output._tracks.some(t => t.metadata.isRelativeToUnixEpoch);
|
||||
const someNotRelative = this.output._tracks.some(t => !t.metadata.isRelativeToUnixEpoch);
|
||||
if (someRelative && someNotRelative) {
|
||||
@@ -449,6 +465,7 @@ export class HlsMuxer extends Muxer {
|
||||
mediaSequence: 0,
|
||||
done: false,
|
||||
singleFile: null,
|
||||
mutex: new AsyncMutex(),
|
||||
};
|
||||
this.playlists.push(playlist);
|
||||
|
||||
@@ -491,6 +508,8 @@ export class HlsMuxer extends Muxer {
|
||||
: [],
|
||||
});
|
||||
}
|
||||
|
||||
release();
|
||||
}
|
||||
|
||||
async getMimeType(): Promise<string> {
|
||||
@@ -509,17 +528,17 @@ export class HlsMuxer extends Muxer {
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-misused-promises
|
||||
override async onTrackClose(track: OutputTrack) {
|
||||
const release = await this.mutex.acquire();
|
||||
const trackData = this.trackDatas.find(x => x.track === track);
|
||||
if (trackData) {
|
||||
trackData.closed = true;
|
||||
}
|
||||
|
||||
const playlist = this.playlists.find(x => x.tracks.includes(track));
|
||||
assert(playlist); // If there isn't one then the assignment algo failed innit
|
||||
|
||||
const release = await playlist.mutex.acquire();
|
||||
|
||||
try {
|
||||
const trackData = this.trackDatas.find(x => x.track === track);
|
||||
if (trackData) {
|
||||
trackData.closed = true;
|
||||
}
|
||||
|
||||
const playlist = this.playlists.find(x => x.tracks.includes(track));
|
||||
assert(playlist); // If there isn't one then the assignment algo failed innit
|
||||
|
||||
await this.advancePlaylist(playlist);
|
||||
} finally {
|
||||
release();
|
||||
@@ -589,18 +608,17 @@ export class HlsMuxer extends Muxer {
|
||||
packet: EncodedPacket,
|
||||
meta?: EncodedVideoChunkMetadata,
|
||||
) {
|
||||
const release = await this.mutex.acquire();
|
||||
const trackData = this.getVideoTrackData(track, meta);
|
||||
const playlist = trackData.playlist;
|
||||
|
||||
const release = await playlist.mutex.acquire();
|
||||
|
||||
try {
|
||||
const trackData = this.getVideoTrackData(track, meta);
|
||||
|
||||
const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key');
|
||||
const adjustedPacket = packet.clone({ timestamp });
|
||||
|
||||
trackData.packets.push(adjustedPacket);
|
||||
|
||||
const playlist = trackData.playlist;
|
||||
|
||||
if (playlist.currentSegmentStartTimestamp === null) {
|
||||
playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp;
|
||||
} else if (!playlist.currentSegmentStartTimestampIsFixed) {
|
||||
@@ -621,18 +639,17 @@ export class HlsMuxer extends Muxer {
|
||||
packet: EncodedPacket,
|
||||
meta?: EncodedAudioChunkMetadata,
|
||||
) {
|
||||
const release = await this.mutex.acquire();
|
||||
const trackData = this.getAudioTrackData(track, meta);
|
||||
const playlist = trackData.playlist;
|
||||
|
||||
const release = await playlist.mutex.acquire();
|
||||
|
||||
try {
|
||||
const trackData = this.getAudioTrackData(track, meta);
|
||||
|
||||
const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key');
|
||||
const adjustedPacket = packet.clone({ timestamp });
|
||||
|
||||
trackData.packets.push(adjustedPacket);
|
||||
|
||||
const playlist = trackData.playlist;
|
||||
|
||||
if (playlist.currentSegmentStartTimestamp === null) {
|
||||
playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp;
|
||||
} else if (!playlist.currentSegmentStartTimestampIsFixed) {
|
||||
@@ -1402,29 +1419,35 @@ export class HlsMuxer extends Muxer {
|
||||
|
||||
this.format._options.onMaster?.(masterPlaylistText);
|
||||
|
||||
let writer: Writer;
|
||||
if (this.numWrittenMasterPlaylists === 0) {
|
||||
// For the first master playlist write, we use the normal root writer getter, so that the target returned by
|
||||
// Output.target emits valid write events.
|
||||
writer = await this.output._getRootWriter();
|
||||
} else {
|
||||
// For subsequent master playlist writes, we *must* obtain a different target in order to overwrite
|
||||
// the file.
|
||||
const target = await this.output._getTarget({
|
||||
path: pathedTarget.rootPath,
|
||||
isRoot: true,
|
||||
mimeType: HLS_MIME_TYPE,
|
||||
});
|
||||
writer = new Writer(target);
|
||||
writer.start();
|
||||
const release = await this.mutex.acquire();
|
||||
|
||||
try {
|
||||
let writer: Writer;
|
||||
if (this.numWrittenMasterPlaylists === 0) {
|
||||
// For the first master playlist write, we use the normal root writer getter, so that the target
|
||||
// returned by Output.target emits valid write events.
|
||||
writer = await this.output._getRootWriter();
|
||||
} else {
|
||||
// For subsequent master playlist writes, we *must* obtain a different target in order to overwrite
|
||||
// the file.
|
||||
const target = await this.output._getTarget({
|
||||
path: pathedTarget.rootPath,
|
||||
isRoot: true,
|
||||
mimeType: HLS_MIME_TYPE,
|
||||
});
|
||||
writer = new Writer(target);
|
||||
writer.start();
|
||||
}
|
||||
|
||||
writer.write(textEncoder.encode(masterPlaylistText));
|
||||
|
||||
await writer.flush();
|
||||
await writer.finalize();
|
||||
|
||||
this.numWrittenMasterPlaylists++;
|
||||
} finally {
|
||||
release();
|
||||
}
|
||||
|
||||
writer.write(textEncoder.encode(masterPlaylistText));
|
||||
|
||||
await writer.flush();
|
||||
await writer.finalize();
|
||||
|
||||
this.numWrittenMasterPlaylists++;
|
||||
}
|
||||
|
||||
private async tryWriteMasterPlaylist() {
|
||||
@@ -1441,8 +1464,8 @@ export class HlsMuxer extends Muxer {
|
||||
}
|
||||
|
||||
async finalize() {
|
||||
assert(this.output._target instanceof PathedTarget);
|
||||
const release = await this.mutex.acquire();
|
||||
const releases = await Promise.all(this.playlists.map(p => p.mutex.acquire()));
|
||||
releases.forEach(release => release());
|
||||
|
||||
for (const trackData of this.trackDatas) {
|
||||
trackData.closed = true;
|
||||
@@ -1455,8 +1478,6 @@ export class HlsMuxer extends Muxer {
|
||||
if (!this.isLive) {
|
||||
await this.writeMasterPlaylist();
|
||||
}
|
||||
|
||||
release();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -130,6 +130,7 @@ export {
|
||||
TargetEvents,
|
||||
TargetRequest,
|
||||
BufferTarget,
|
||||
BufferTargetOptions,
|
||||
FilePathTarget,
|
||||
FilePathTargetOptions,
|
||||
NullTarget,
|
||||
@@ -141,6 +142,7 @@ export {
|
||||
} from './target';
|
||||
export {
|
||||
AnyIterable,
|
||||
ConcurrentRunner,
|
||||
EventEmitter,
|
||||
EventListenerOptions,
|
||||
FilePath,
|
||||
|
||||
+34
-24
@@ -267,6 +267,8 @@ class VideoEncoderWrapper {
|
||||
*/
|
||||
private error: Error | null = null;
|
||||
|
||||
private lastMuxerPromise: Promise<void> = Promise.resolve();
|
||||
|
||||
constructor(private source: VideoSource, private encodingConfig: VideoEncodingConfig) {
|
||||
const sizeChangeBehavior = encodingConfig.sizeChangeBehavior ?? 'deny';
|
||||
if (['fill', 'contain', 'cover'].includes(sizeChangeBehavior) && encodingConfig.transform?.fit !== undefined) {
|
||||
@@ -648,7 +650,7 @@ class VideoEncoderWrapper {
|
||||
}
|
||||
}
|
||||
|
||||
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
|
||||
await this.lastMuxerPromise; // Allow the writer to apply backpressure
|
||||
}
|
||||
} finally {
|
||||
for (const sample of samplesToEncode) {
|
||||
@@ -711,10 +713,11 @@ class VideoEncoderWrapper {
|
||||
maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet);
|
||||
|
||||
this.encodingConfig.onEncodedPacket?.(packet, meta);
|
||||
void this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
this.lastMuxerPromise
|
||||
= this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
};
|
||||
|
||||
await this.customEncoder.init();
|
||||
@@ -802,10 +805,11 @@ class VideoEncoderWrapper {
|
||||
maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet);
|
||||
|
||||
this.encodingConfig.onEncodedPacket?.(packet, meta);
|
||||
void this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
this.lastMuxerPromise
|
||||
= this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
};
|
||||
|
||||
const stack = new Error('Encoding error').stack;
|
||||
@@ -1787,6 +1791,7 @@ class AudioEncoderWrapper {
|
||||
* So, we keep track of the encoder error and throw it as soon as we get the chance.
|
||||
*/
|
||||
private error: Error | null = null;
|
||||
private lastMuxerPromise: Promise<void> = Promise.resolve();
|
||||
|
||||
constructor(private source: AudioSource, private encodingConfig: AudioEncodingConfig) {}
|
||||
|
||||
@@ -1950,7 +1955,7 @@ class AudioEncoderWrapper {
|
||||
await promise;
|
||||
}
|
||||
|
||||
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
|
||||
await this.lastMuxerPromise; // Allow the writer to apply backpressure
|
||||
} else if (this.isPcmEncoder) {
|
||||
await this.doPcmEncoding(audioSample, shouldClose);
|
||||
} else {
|
||||
@@ -1967,7 +1972,7 @@ class AudioEncoderWrapper {
|
||||
await new Promise(resolve => this.encoder!.addEventListener('dequeue', resolve, { once: true }));
|
||||
}
|
||||
|
||||
await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure
|
||||
await this.lastMuxerPromise; // Allow the writer to apply backpressure
|
||||
}
|
||||
} finally {
|
||||
if (shouldClose) {
|
||||
@@ -2080,10 +2085,11 @@ class AudioEncoderWrapper {
|
||||
}
|
||||
|
||||
this.encodingConfig.onEncodedPacket?.(packet, meta);
|
||||
void this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
this.lastMuxerPromise
|
||||
= this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
};
|
||||
|
||||
await this.customEncoder.init();
|
||||
@@ -2145,10 +2151,11 @@ class AudioEncoderWrapper {
|
||||
});
|
||||
|
||||
this.encodingConfig.onEncodedPacket?.(packet, meta);
|
||||
void this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
this.lastMuxerPromise
|
||||
= this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta)
|
||||
.catch((error) => {
|
||||
this.error ??= error;
|
||||
});
|
||||
},
|
||||
error: (error) => {
|
||||
error.stack = stack; // Provide a more useful stack trace
|
||||
@@ -2803,6 +2810,8 @@ export class TextSubtitleSource extends SubtitleSource {
|
||||
private _parser: SubtitleParser;
|
||||
/** @internal */
|
||||
private _error: Error | null = null;
|
||||
/** @internal */
|
||||
private _lastMuxerPromise: Promise<void> = Promise.resolve();
|
||||
|
||||
/** Creates a new {@link TextSubtitleSource} where added text chunks are in the specified `codec`. */
|
||||
constructor(codec: SubtitleCodec) {
|
||||
@@ -2811,10 +2820,11 @@ export class TextSubtitleSource extends SubtitleSource {
|
||||
this._parser = new SubtitleParser({
|
||||
codec,
|
||||
output: (cue, metadata) => {
|
||||
void this._connectedTrack?.output._muxer.addSubtitleCue(this._connectedTrack, cue, metadata)
|
||||
.catch((error) => {
|
||||
this._error ??= error;
|
||||
});
|
||||
this._lastMuxerPromise
|
||||
= this._connectedTrack!.output._muxer.addSubtitleCue(this._connectedTrack!, cue, metadata)
|
||||
.catch((error) => {
|
||||
this._error ??= error;
|
||||
});
|
||||
},
|
||||
});
|
||||
}
|
||||
@@ -2836,7 +2846,7 @@ export class TextSubtitleSource extends SubtitleSource {
|
||||
this._ensureValidAdd();
|
||||
this._parser.parse(text);
|
||||
|
||||
return this._connectedTrack!.output._muxer.mutex.currentPromise;
|
||||
return this._lastMuxerPromise; // Allow the writer to apply backpressure
|
||||
}
|
||||
|
||||
/** @internal */
|
||||
|
||||
+63
@@ -1235,3 +1235,66 @@ export class EventEmitter<TEvents extends Record<string, unknown>> {
|
||||
}
|
||||
|
||||
export const ceilToMultipleOfTwo = (value: number) => Math.ceil(value / 2) * 2;
|
||||
|
||||
/**
|
||||
* Utility class for running async functions in parallel up to a certain level of parallelism. Can be used to apply
|
||||
* backpressure only if the concurrency level would be exceeded.
|
||||
*
|
||||
* @group Miscellaneous
|
||||
* @public
|
||||
*/
|
||||
export class ConcurrentRunner {
|
||||
/** @internal */
|
||||
_queue: Promise<unknown>[] = [];
|
||||
/** @internal */
|
||||
_errored = false;
|
||||
|
||||
/**
|
||||
* The maximum number of in-flight promises. You can also think of it as the "high water mark".
|
||||
* You can set this value to dynamically change the level of parallelism.
|
||||
*/
|
||||
parallelism: number;
|
||||
|
||||
constructor(parallelism: number) {
|
||||
this.parallelism = parallelism;
|
||||
}
|
||||
|
||||
/** Whether any function has errored. The runner is effectively bricked if this is `true`, by design. */
|
||||
get errored() {
|
||||
return this._errored;
|
||||
}
|
||||
|
||||
/** The number of tasks currently running. */
|
||||
get inFlightCount() {
|
||||
return this._queue.length;
|
||||
}
|
||||
|
||||
/**
|
||||
* Schedules an async function to be run. If the maximum allowed level of parallelism has not yet been reached,
|
||||
* the function will be executed immediately and `run()` will resolve immediately. Otherwise, the function will be
|
||||
* called as soon as any currently-running function finishes, and `run()` will only resolve then.
|
||||
*
|
||||
* Throws if the runner is errored.
|
||||
*/
|
||||
async run(fn: () => Promise<unknown>) {
|
||||
if (this._errored) {
|
||||
await Promise.race(this._queue); // Will surface the error
|
||||
}
|
||||
|
||||
while (this._queue.length >= this.parallelism) {
|
||||
await Promise.race(this._queue);
|
||||
}
|
||||
|
||||
const promise = fn();
|
||||
this._queue.push(promise);
|
||||
|
||||
void promise
|
||||
.then(() => removeItem(this._queue, promise))
|
||||
.catch(() => this._errored = true);
|
||||
}
|
||||
|
||||
/** Waits for all currently running functions to finish. Throws if the runner is errored. */
|
||||
async flush() {
|
||||
await Promise.all(this._queue);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -324,6 +324,11 @@ export type OutputOptions<
|
||||
* When this is a function, it will only be called if an init target is needed.
|
||||
*/
|
||||
initTarget?: T | (() => MaybePromise<T>);
|
||||
/**
|
||||
* Optional; a callback to be called at the end of {@link Output.finalize}. Can be used to run logic once the
|
||||
* output has completed. If a promise is returned, it will be awaited internally by {@link Output.finalize}.
|
||||
*/
|
||||
onFinalize?: () => MaybePromise<unknown>;
|
||||
};
|
||||
|
||||
/**
|
||||
@@ -369,6 +374,8 @@ export class Output<
|
||||
/** @internal */
|
||||
private _initTarget: T | (() => MaybePromise<T>) | null;
|
||||
/** @internal */
|
||||
_onFinalize: (() => MaybePromise<unknown>) | null = null;
|
||||
/** @internal */
|
||||
_muxer: Muxer;
|
||||
/** @internal */
|
||||
_targets = new Set<Target>();
|
||||
@@ -449,9 +456,13 @@ export class Output<
|
||||
+ ' a Target.',
|
||||
);
|
||||
}
|
||||
if (options.onFinalize !== undefined && typeof options.onFinalize !== 'function') {
|
||||
throw new TypeError('options.onFinalize, when provided, must be a function.');
|
||||
}
|
||||
|
||||
this.format = options.format;
|
||||
this._target = options.target;
|
||||
this._onFinalize = options.onFinalize ?? null;
|
||||
|
||||
this._initTarget = options.initTarget ?? null;
|
||||
if (this._initTarget instanceof Target) {
|
||||
@@ -881,6 +892,10 @@ export class Output<
|
||||
}
|
||||
}
|
||||
|
||||
if (this._onFinalize) {
|
||||
await this._onFinalize();
|
||||
}
|
||||
|
||||
this.state = 'finalized';
|
||||
} finally {
|
||||
release();
|
||||
|
||||
+17
-2
@@ -70,6 +70,8 @@ export abstract class SegmentedInput {
|
||||
firstSegment: Segment | null = null;
|
||||
firstSegmentFirstTimestamps = new WeakMap<Segment, number>();
|
||||
|
||||
firstTimestampCache = new WeakMap<Input, number>();
|
||||
|
||||
constructor(input: Input, path: string, trackDeclarations: SegmentedInputTrackDeclaration[] | null) {
|
||||
this.input = input;
|
||||
this.path = path;
|
||||
@@ -152,6 +154,19 @@ export abstract class SegmentedInput {
|
||||
})();
|
||||
}
|
||||
|
||||
// This operation is done a lot and can be semi-expensive, so it's good to have a cache for it
|
||||
async getFirstTimestampForInput(input: Input) {
|
||||
const existing = this.firstTimestampCache.get(input);
|
||||
if (existing !== undefined) {
|
||||
return existing;
|
||||
}
|
||||
|
||||
const firstTimestamp = await input.getFirstTimestamp();
|
||||
this.firstTimestampCache.set(input, firstTimestamp);
|
||||
|
||||
return firstTimestamp;
|
||||
}
|
||||
|
||||
async getMediaOffset(segment: Segment, input: Input) {
|
||||
const firstSegment = segment.firstSegment ?? segment;
|
||||
|
||||
@@ -160,7 +175,7 @@ export abstract class SegmentedInput {
|
||||
firstSegmentFirstTimestamp = this.firstSegmentFirstTimestamps.get(firstSegment)!;
|
||||
} else {
|
||||
const firstInput = this.getInputForSegment(firstSegment);
|
||||
firstSegmentFirstTimestamp = await firstInput.getFirstTimestamp();
|
||||
firstSegmentFirstTimestamp = await this.getFirstTimestampForInput(firstInput);
|
||||
this.firstSegmentFirstTimestamps.set(firstSegment, firstSegmentFirstTimestamp);
|
||||
}
|
||||
|
||||
@@ -168,7 +183,7 @@ export abstract class SegmentedInput {
|
||||
return firstSegment.timestamp - firstSegmentFirstTimestamp;
|
||||
}
|
||||
|
||||
const segmentFirstTimestamp = await input.getFirstTimestamp();
|
||||
const segmentFirstTimestamp = await this.getFirstTimestampForInput(input);
|
||||
const segmentElapsed = segment.timestamp - firstSegment.timestamp;
|
||||
const inputElapsed = segmentFirstTimestamp - firstSegmentFirstTimestamp;
|
||||
const difference = inputElapsed - segmentElapsed;
|
||||
|
||||
+33
-1
@@ -89,6 +89,22 @@ export abstract class Target extends EventEmitter<TargetEvents> {
|
||||
const ARRAY_BUFFER_INITIAL_SIZE = 2 ** 16;
|
||||
const ARRAY_BUFFER_MAX_SIZE = 2 ** 32;
|
||||
|
||||
/**
|
||||
* Options for {@link BufferTarget}.
|
||||
* @group Output targets
|
||||
* @public
|
||||
*/
|
||||
export type BufferTargetOptions = {
|
||||
/**
|
||||
* Called once the target has been finalized, with the complete output buffer. If you return a promise, it will be
|
||||
* used to apply backpressure internally.
|
||||
*
|
||||
* One use for this callback is for uploading to a server where the full buffer must be known before
|
||||
* sending (e.g. S3 PutObject) and stream-uploading is not an option.
|
||||
*/
|
||||
onFinalize?: (buffer: ArrayBuffer) => MaybePromise<unknown>;
|
||||
};
|
||||
|
||||
/**
|
||||
* A target that writes data directly into an ArrayBuffer in memory. Great for performance, but not suitable for very
|
||||
* large files. The buffer will be available once the output has been finalized.
|
||||
@@ -107,11 +123,22 @@ export class BufferTarget extends Target {
|
||||
_maxPos = 0;
|
||||
/** @internal */
|
||||
_supportsResize: boolean;
|
||||
/** @internal */
|
||||
_options: BufferTargetOptions;
|
||||
|
||||
/** Creates a new {@link BufferTarget}. The buffer holding the data will be created and managed internally. */
|
||||
constructor() {
|
||||
constructor(options: BufferTargetOptions = {}) {
|
||||
super();
|
||||
|
||||
if (!options || typeof options !== 'object') {
|
||||
throw new TypeError('BufferTarget options, when provided, must be an object.');
|
||||
}
|
||||
if (options.onFinalize !== undefined && typeof options.onFinalize !== 'function') {
|
||||
throw new TypeError('options.onFinalize, when provided, must be a function.');
|
||||
}
|
||||
|
||||
this._options = options;
|
||||
|
||||
this._supportsResize = 'resize' in new ArrayBuffer(0);
|
||||
if (this._supportsResize) {
|
||||
try {
|
||||
@@ -178,6 +205,11 @@ export class BufferTarget extends Target {
|
||||
/** @internal */
|
||||
async _finalize() {
|
||||
this.buffer = this._buffer.slice(0, this._maxPos);
|
||||
|
||||
if (this._options.onFinalize) {
|
||||
await this._options.onFinalize(this.buffer);
|
||||
}
|
||||
|
||||
this._emit('finalized');
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,289 @@
|
||||
import { expect, test } from 'vitest';
|
||||
import { ConcurrentRunner, promiseWithResolvers } from '../../src/misc.js';
|
||||
|
||||
test('parallelism of 1 runs tasks strictly sequentially', async () => {
|
||||
const runner = new ConcurrentRunner(1);
|
||||
let entered = false;
|
||||
let overlapped = false;
|
||||
|
||||
const makeTask = () => async () => {
|
||||
if (entered) {
|
||||
overlapped = true;
|
||||
}
|
||||
entered = true;
|
||||
await Promise.resolve();
|
||||
await Promise.resolve();
|
||||
entered = false;
|
||||
};
|
||||
|
||||
const runPromises: Promise<void>[] = [];
|
||||
for (let i = 0; i < 5; i++) {
|
||||
runPromises.push(runner.run(makeTask()));
|
||||
}
|
||||
|
||||
await Promise.all(runPromises);
|
||||
await runner.flush();
|
||||
|
||||
expect(overlapped).toBe(false);
|
||||
expect(entered).toBe(false);
|
||||
});
|
||||
|
||||
test('run schedules tasks up to parallelism without waiting', async () => {
|
||||
const runner = new ConcurrentRunner(3);
|
||||
const d1 = promiseWithResolvers();
|
||||
const d2 = promiseWithResolvers();
|
||||
const d3 = promiseWithResolvers();
|
||||
|
||||
let started = 0;
|
||||
|
||||
await runner.run(async () => {
|
||||
started++;
|
||||
await d1.promise;
|
||||
});
|
||||
await runner.run(async () => {
|
||||
started++;
|
||||
await d2.promise;
|
||||
});
|
||||
await runner.run(async () => {
|
||||
started++;
|
||||
await d3.promise;
|
||||
});
|
||||
|
||||
// All three should have started synchronously since we're under parallelism
|
||||
expect(started).toBe(3);
|
||||
|
||||
d1.resolve();
|
||||
d2.resolve();
|
||||
d3.resolve();
|
||||
await runner.flush();
|
||||
});
|
||||
|
||||
test('run blocks when the queue is full and resumes as slots free up', async () => {
|
||||
const runner = new ConcurrentRunner(2);
|
||||
const d1 = promiseWithResolvers();
|
||||
const d2 = promiseWithResolvers();
|
||||
const d3 = promiseWithResolvers();
|
||||
|
||||
let startedThird = false;
|
||||
|
||||
await runner.run(() => d1.promise);
|
||||
await runner.run(() => d2.promise);
|
||||
|
||||
// Kick off a third run. It must not resolve until one of the first two finishes.
|
||||
const thirdRun = runner.run(async () => {
|
||||
startedThird = true;
|
||||
await d3.promise;
|
||||
});
|
||||
|
||||
await flushMicrotasks();
|
||||
expect(startedThird).toBe(false);
|
||||
|
||||
// Free a slot by resolving the first task.
|
||||
d1.resolve();
|
||||
await thirdRun;
|
||||
|
||||
expect(startedThird).toBe(true);
|
||||
|
||||
d2.resolve();
|
||||
d3.resolve();
|
||||
await runner.flush();
|
||||
});
|
||||
|
||||
test('simultaneous run calls respect parallelism', async () => {
|
||||
const runner = new ConcurrentRunner(2);
|
||||
const deferreds = [
|
||||
promiseWithResolvers(),
|
||||
promiseWithResolvers(),
|
||||
promiseWithResolvers(),
|
||||
promiseWithResolvers(),
|
||||
];
|
||||
const started: number[] = [];
|
||||
|
||||
const runPromises = deferreds.map((d, i) => runner.run(async () => {
|
||||
started.push(i);
|
||||
await d.promise;
|
||||
}));
|
||||
|
||||
// Give the runner a chance to start the first batch, then verify only the first two ran.
|
||||
await flushMicrotasks();
|
||||
expect(started).toEqual([0, 1]);
|
||||
|
||||
// The first two run() calls should have resolved already (they found open slots),
|
||||
// the last two should still be waiting.
|
||||
let settled = 0;
|
||||
void Promise.all(runPromises).then(() => settled++);
|
||||
await flushMicrotasks();
|
||||
expect(settled).toBe(0);
|
||||
|
||||
// Unblock the first task. That should let task 2 in.
|
||||
deferreds[0]!.resolve();
|
||||
await flushMicrotasks();
|
||||
expect(started).toEqual([0, 1, 2]);
|
||||
|
||||
// Unblock the second task. That should let task 3 in.
|
||||
deferreds[1]!.resolve();
|
||||
await flushMicrotasks();
|
||||
expect(started).toEqual([0, 1, 2, 3]);
|
||||
|
||||
// All four run() calls should resolve once their slot has been acquired.
|
||||
await Promise.all(runPromises);
|
||||
|
||||
deferreds[2]!.resolve();
|
||||
deferreds[3]!.resolve();
|
||||
await runner.flush();
|
||||
});
|
||||
|
||||
test('errored task surfaces on the next run call', async () => {
|
||||
const runner = new ConcurrentRunner(2);
|
||||
const error = new Error('boom');
|
||||
|
||||
expect(runner.errored).toBe(false);
|
||||
|
||||
await runner.run(async () => {
|
||||
throw error;
|
||||
});
|
||||
|
||||
await flushMicrotasks();
|
||||
expect(runner.errored).toBe(true);
|
||||
|
||||
await expect(runner.run(async () => {})).rejects.toBe(error);
|
||||
});
|
||||
|
||||
test('errored task surfaces on flush', async () => {
|
||||
const runner = new ConcurrentRunner(2);
|
||||
const error = new Error('kaboom');
|
||||
|
||||
await runner.run(async () => {
|
||||
throw error;
|
||||
});
|
||||
|
||||
await expect(runner.flush()).rejects.toBe(error);
|
||||
expect(runner.errored).toBe(true);
|
||||
});
|
||||
|
||||
test('error from a slow task surfaces on subsequent run even after other tasks completed', async () => {
|
||||
const runner = new ConcurrentRunner(2);
|
||||
const slow = promiseWithResolvers();
|
||||
const error = new Error('late');
|
||||
|
||||
await runner.run(async () => {
|
||||
await slow.promise;
|
||||
throw error;
|
||||
});
|
||||
await runner.run(async () => {});
|
||||
|
||||
// Let the fast task finish cleanly.
|
||||
await flushMicrotasks();
|
||||
expect(runner.errored).toBe(false);
|
||||
|
||||
// Now let the slow task reject.
|
||||
slow.resolve();
|
||||
await flushMicrotasks();
|
||||
expect(runner.errored).toBe(true);
|
||||
|
||||
await expect(runner.run(async () => {})).rejects.toBe(error);
|
||||
});
|
||||
|
||||
test('parallelism can be mutated at runtime to grow or shrink the in-flight limit', async () => {
|
||||
const runner = new ConcurrentRunner(1);
|
||||
const d1 = promiseWithResolvers();
|
||||
const d2 = promiseWithResolvers();
|
||||
const d3 = promiseWithResolvers();
|
||||
|
||||
let startedSecond = false;
|
||||
let startedThird = false;
|
||||
|
||||
await runner.run(() => d1.promise);
|
||||
expect(runner.inFlightCount).toBe(1);
|
||||
|
||||
// Grow to 2 before scheduling the next task — second run sees the new value and starts immediately.
|
||||
runner.parallelism = 2;
|
||||
await runner.run(async () => {
|
||||
startedSecond = true;
|
||||
await d2.promise;
|
||||
});
|
||||
expect(startedSecond).toBe(true);
|
||||
expect(runner.inFlightCount).toBe(2);
|
||||
|
||||
// Shrink to 1 while 2 are in flight. No task is cancelled; a new run() must wait until queue < 1.
|
||||
runner.parallelism = 1;
|
||||
const third = runner.run(async () => {
|
||||
startedThird = true;
|
||||
await d3.promise;
|
||||
});
|
||||
await flushMicrotasks();
|
||||
expect(startedThird).toBe(false);
|
||||
|
||||
// Draining one frees a slot but queue is still at 1 (>= new parallelism), third stays blocked.
|
||||
d1.resolve();
|
||||
await flushMicrotasks();
|
||||
expect(startedThird).toBe(false);
|
||||
|
||||
// Draining the second lets third in.
|
||||
d2.resolve();
|
||||
await third;
|
||||
expect(startedThird).toBe(true);
|
||||
|
||||
d3.resolve();
|
||||
await runner.flush();
|
||||
});
|
||||
|
||||
test('inFlightCount tracks currently running tasks', async () => {
|
||||
const runner = new ConcurrentRunner(3);
|
||||
const d1 = promiseWithResolvers();
|
||||
const d2 = promiseWithResolvers();
|
||||
|
||||
expect(runner.inFlightCount).toBe(0);
|
||||
|
||||
await runner.run(() => d1.promise);
|
||||
expect(runner.inFlightCount).toBe(1);
|
||||
|
||||
await runner.run(() => d2.promise);
|
||||
expect(runner.inFlightCount).toBe(2);
|
||||
|
||||
d1.resolve();
|
||||
await flushMicrotasks();
|
||||
expect(runner.inFlightCount).toBe(1);
|
||||
|
||||
d2.resolve();
|
||||
await runner.flush();
|
||||
expect(runner.inFlightCount).toBe(0);
|
||||
});
|
||||
|
||||
test('flush waits for all in-flight tasks', async () => {
|
||||
const runner = new ConcurrentRunner(3);
|
||||
const deferreds = [promiseWithResolvers(), promiseWithResolvers(), promiseWithResolvers()];
|
||||
const completed: number[] = [];
|
||||
|
||||
for (let i = 0; i < deferreds.length; i++) {
|
||||
await runner.run(async () => {
|
||||
await deferreds[i]!.promise;
|
||||
completed.push(i);
|
||||
});
|
||||
}
|
||||
|
||||
let flushResolved = false;
|
||||
const flushPromise = runner.flush().then(() => {
|
||||
flushResolved = true;
|
||||
});
|
||||
|
||||
await flushMicrotasks();
|
||||
expect(flushResolved).toBe(false);
|
||||
|
||||
deferreds[0]!.resolve();
|
||||
await flushMicrotasks();
|
||||
expect(flushResolved).toBe(false);
|
||||
|
||||
deferreds[1]!.resolve();
|
||||
deferreds[2]!.resolve();
|
||||
await flushPromise;
|
||||
|
||||
expect(flushResolved).toBe(true);
|
||||
expect(completed.sort()).toEqual([0, 1, 2]);
|
||||
});
|
||||
|
||||
const flushMicrotasks = async (iterations = 10) => {
|
||||
for (let i = 0; i < iterations; i++) {
|
||||
await Promise.resolve();
|
||||
}
|
||||
};
|
||||
@@ -0,0 +1,36 @@
|
||||
import { expect, test } from 'vitest';
|
||||
import { Output } from '../../src/output.js';
|
||||
import { MkvOutputFormat } from '../../src/output-format.js';
|
||||
import { BufferTarget } from '../../src/target.js';
|
||||
import { EncodedVideoPacketSource } from '../../src/media-source.js';
|
||||
import { EncodedPacket } from '../../src/packet.js';
|
||||
|
||||
test('Output, onFinalize', async () => {
|
||||
let callCount = 0;
|
||||
const output = new Output({
|
||||
format: new MkvOutputFormat(),
|
||||
target: new BufferTarget(),
|
||||
onFinalize: async () => {
|
||||
await new Promise(resolve => setTimeout(resolve, 200));
|
||||
|
||||
callCount++;
|
||||
},
|
||||
});
|
||||
|
||||
const source = new EncodedVideoPacketSource('avc');
|
||||
output.addVideoTrack(source);
|
||||
|
||||
await output.start();
|
||||
|
||||
await source.add(new EncodedPacket(new Uint8Array(1024), 'key', 0, 0.5), {
|
||||
decoderConfig: {
|
||||
codec: 'avc1.640028',
|
||||
codedWidth: 1920,
|
||||
codedHeight: 1080,
|
||||
},
|
||||
});
|
||||
|
||||
await output.finalize();
|
||||
|
||||
expect(callCount).toBe(1);
|
||||
});
|
||||
@@ -0,0 +1,42 @@
|
||||
import path from 'node:path';
|
||||
import { expect, test } from 'vitest';
|
||||
import { Input } from '../../src/input.js';
|
||||
import { FilePathSource } from '../../src/source.js';
|
||||
import { ALL_FORMATS } from '../../src/input-format.js';
|
||||
import { Output } from '../../src/output.js';
|
||||
import { BufferTarget } from '../../src/target.js';
|
||||
import { Mp4OutputFormat } from '../../src/output-format.js';
|
||||
import { Conversion } from '../../src/conversion.js';
|
||||
|
||||
const __dirname = new URL('.', import.meta.url).pathname;
|
||||
|
||||
const samplePath = path.join(__dirname, '../public/video.mp4');
|
||||
|
||||
test('BufferTarget onFinalize callback', async () => {
|
||||
let received: ArrayBuffer | null = null;
|
||||
let asyncCallbackDone = false;
|
||||
|
||||
using input = new Input({
|
||||
source: new FilePathSource(samplePath),
|
||||
formats: ALL_FORMATS,
|
||||
});
|
||||
|
||||
const output = new Output({
|
||||
format: new Mp4OutputFormat(),
|
||||
target: new BufferTarget({
|
||||
onFinalize: async (buffer) => {
|
||||
received = buffer;
|
||||
await new Promise(resolve => setTimeout(resolve, 20));
|
||||
asyncCallbackDone = true;
|
||||
},
|
||||
}),
|
||||
});
|
||||
|
||||
const conversion = await Conversion.init({ input, output, showWarnings: false });
|
||||
await conversion.execute();
|
||||
|
||||
expect(received).not.toBeNull();
|
||||
expect(received).toBe(output.target.buffer);
|
||||
expect(asyncCallbackDone).toBe(true);
|
||||
expect(output.target.buffer!.byteLength).toBeGreaterThan(0);
|
||||
});
|
||||
@@ -0,0 +1,10 @@
|
||||
MENTION APPEND-ONLY IN UPLOAD EXAMPLE IN WRITING HLS
|
||||
|
||||
also I wish there was a more explicit way this was enforced and would error at runtime.
|
||||
writablestream target?
|
||||
|
||||
Thoughts:
|
||||
So, a certain "lookahead" logic is definitely needed. The question is if this is a per-demuxer thing or a general thing instead. The demuxer could get in a "packet query" that specifies things like "I am interested in the next 20 seconds guaranteed", allowing the demuxer to pre-fetch more intelligently. The alternative would be some sort of demuxer-agnostic approach where there is a magical "packet requester" that has to be segment-aware. I'm actually not sure if that's good.
|
||||
|
||||
|
||||
- Remove createInputFrom. Much more elegant solution now that PathedSource is a thing!!
|
||||
Reference in New Issue
Block a user