Compare commits

..
Author SHA1 Message Date
Vanilagy 9de93877bd Bump to beta 6 2026-04-23 16:14:24 +02:00
Vanilagy 8809bb311b Remove this 2026-04-23 14:47:48 +02:00
b05cdbe7e0 fix: two long-running HLS transcode issues (#355)
* fix: release targets from Output._targets on finalize

Long-running HLS transcodes leak memory. Every finalized BufferTarget
stays in _targets until the outer Output closes, pinning its buffer.
Writer.finalize() already does this cleanup for writer-based flows;
extend it to buffer-finalize paths via the public 'finalized' event.

* fix: ReadOrchestrator LRU eviction picks only drained workers

assert(pendingSlices.length === 0) fires under heavy concurrent reads
(e.g. multi-rendition HLS decode from BlobSource). LRU filter only
checked !running; workers with queued slices could be evicted.
Add pendingSlices.length === 0 to the filter.

* fix: export AppendOnlyStreamTarget from index

Missing from the re-export; public docs import it by name.

* fix: join HLS init segment path with root + playlist path

Init path went through _getTarget bare; segments got joined with rootPath + playlist.path. Playlist-relative URI then can't resolve when the playlist lives in a subdirectory.

* Fix targets not being cleaned up, fix paused workers with remaining pending slices, modify doc block, fixed isRoot not being changed on proxied requests

---------

Co-authored-by: Vanilagy <[email protected]>
2026-04-23 14:37:56 +02:00
Vanilagy 7acae8dda1 Add a track disposition section to Writing HLS 2026-04-23 11:06:50 +02:00
18 changed files with 185 additions and 318 deletions
-198
View File
@@ -1,198 +0,0 @@
================================================================================
MEDIABUNNY HLS BRANCH - BUG FIXES AND BEHAVIOR CHANGES TO EXISTING CODE
For use in release notes. Only covers changes to pre-existing code.
New HLS/CMAF features are not listed here.
================================================================================
============================================================
BUG FIXES
============================================================
1. SourceRef race conditions (commit aa810fc)
- Source cache entry was added BEFORE the reference was fully created,
allowing premature garbage collection of cached sources
- Cache eviction variable pointed to wrong ref (local instead of entry)
- `count > MAX_SOURCE_CACHE_SIZE` off-by-one -> `count >= MAX_SOURCE_CACHE_SIZE`
- Input.source and Input.target getters could return unused/stale instances
when using async callbacks; now properly tracked via _getRootSourceRef()
2. Track closing race conditions in all muxers (commit 94678dd)
- ISOBMFF, Matroska, MPEG-TS, and OGG muxers used `trackData.track.source._closed`
to check if tracks were closed. This was racy because the closed state could
change between check and use. Added explicit `closed: boolean` field to all
muxer track data structures, set synchronously in onTrackClose().
3. Target.slice() validation bug (commit 7b70756)
- Validation used `!Number.isInteger(offset) && offset < 0` (AND), meaning
negative non-integer offsets slipped through. Fixed to use OR (`||`).
4. MPEG-TS demuxer: off-by-one errors in reorder buffer logic
- Rewind loop changed from `reorderSize` to `reorderSize + 1` iterations,
fixing packet duration calculation near seek points
- End-of-stream rewind loop changed from `reorderSize - 1` to `reorderSize`
- Flush condition changed from `>= reorderSize` to `> reorderSize`, keeping
one extra packet in the buffer before flushing for correct presentation
order computation
5. MPEG-TS demuxer: first chunk assertion crash (line 1031-1035)
- Old code asserted that a key frame is always found in the first chunk,
crashing on files where that assumption doesn't hold (e.g., HLS segments
starting mid-GOP). Changed to gracefully return null.
6. MPEG-TS demuxer: video parameters only searched in first packet
- Some muxers place SPS/PPS in later packets. The demuxer now loops through
multiple packets to find AVC/HEVC decoder configuration records instead of
only checking the first one.
7. ISOBMFF: duration calculation didn't account for sample duration
- `lastPresentedSample()` found the sample with the highest timestamp but
didn't add its duration. Replaced with `presentationSpan()` that correctly
computes `maxEndTimestamp - minTimestamp`.
8. ISOBMFF demuxer: assertion crash on unavailable data
- Demuxer asserted that read slices are non-null. Changed to gracefully
return null when data is outside available range (important for streaming/
progressive scenarios).
9. ISOBMFF demuxer: multiple trun boxes per fragment rejected
- Old code logged a warning and skipped the second trun box. Now correctly
accumulates samples across multiple trun boxes per track per fragment,
maintaining cumulative offset and timestamp state.
10. Conversion API: track count validation before fan-out consideration
(commit 2f576aa)
- Tracks were validated against the max count before considering that some
track options had `discard: true`. Now validates inside the fan-out loop.
11. MP3 muxer: frame positions recorded at wrong offset
- Frame byte positions were recorded AFTER writing packet data. Now recorded
BEFORE writing, which is correct for Xing TOC frame offset calculations.
12. FLAC demuxer: missing STREAMINFO validation
- Added explicit error when STREAMINFO metadata block is missing,
producing a clear "Corrupted FLAC file" message instead of undefined
behavior downstream.
============================================================
BEHAVIOR CHANGES
============================================================
1. Default keyFrameInterval changed from 5 seconds to 2 seconds
(commit c6e505a)
Aligns with default HLS segment duration. Affects all video encoding
that doesn't explicitly set keyFrameInterval.
2. Conversion API copies input pairability graph by default (commit 6410ea9)
When track group is not explicitly specified, conversion now auto-creates
OutputTrackGroups that mirror the input track pairing relationships. This
means converted outputs preserve which tracks are meant to play together.
3. Matroska demuxer: removed implicit track sorting by default disposition
Tracks are no longer sorted so that default=true tracks come first.
Now all tracks have `primary: false` by default. Callers should use the
getPrimaryVideoTrack()/getPrimaryAudioTrack() methods for selection.
4. Removed superfluous end-position seeks in muxer finalize methods
(commit 7b70756)
MP3, FLAC, and Matroska muxers no longer seek to the end of file after
writing final metadata. The seek was unnecessary since the writer is
done at that point.
5. clampCropRectangle now returns a new object instead of mutating input
(src/sample.ts:1123)
6. Date.now() -> performance.now() in FinalizationRegistry callback
(src/sample.ts:48)
Uses monotonic time for timing measurements.
7. Codec string parsing: mp4a.40.34 now correctly identified as MP3
(src/codec.ts:677-686)
Previously would have matched the aac prefix check. MP3 check now
runs first and includes this codec string.
8. UrlSource: servers without Content-Length now supported
Previously, if the server returned 200 without a Content-Length header,
UrlSource threw an error. Now it gracefully handles this by downloading
the entire resource with an unbounded worker (targetPos = Infinity,
strictTarget = false), determining file size once the stream ends.
9. UrlSource: file size probing request removed
The old UrlSource always made a dedicated initial `Range: bytes=0-`
request solely to probe file size and range request support. File size
is now determined lazily from the response headers of the first actual
read, removing the extra round-trip.
10. UrlSource: non-range-request server warnings deduplicated per origin
The "server did not respond with 206 Partial Content" warning is now
emitted at most once per origin instead of on every request.
11. ReadableStreamSource: reader now canceled on dispose
ReadableStreamSource._dispose() now calls `this._reader?.cancel()`,
properly releasing the underlying stream resource.
12. ReadOrchestrator: worker queue system
When the maximum worker count is reached, reads are now queued and
dispatched when a worker becomes free, instead of evicting running
workers which could abort in-flight fetches.
============================================================
DEPRECATIONS
============================================================
1. InputTrack sync property getters (commit 051578c)
All sync getters on InputTrack/InputVideoTrack/InputAudioTrack are now
deprecated in favor of async methods:
.codec -> await .getCodec()
.languageCode -> await .getLanguageCode()
.name -> await .getName()
.timeResolution -> await .getTimeResolution()
.disposition -> await .getDisposition()
.displayWidth -> await .getDisplayWidth()
.displayHeight -> await .getDisplayHeight()
.rotation -> await .getRotation()
(etc.)
The sync getters still work for non-HLS inputs but throw when the
backing requires async resolution (e.g., HLS tracks before hydration).
2. Source.onread callback -> source.on('read', handler)
The onread setter still works but is marked @deprecated.
3. Target.onwrite callback -> target.on('write', handler)
Same as above.
============================================================
REMOVALS
============================================================
1. Target.onfinalized callback (commit 6410ea9)
Removed entirely. Use the 'finalized' event: target.on('finalized', ...).
2. InputTrackDescriptor concept (commit 051578c)
The entire InputTrackDescriptor API was removed. Its functionality
(pairable tracks, primary track selection) is now directly on InputTrack.
3. InputTrack[] option for ConversionOptions.tracks (commit 606ee87)
ConversionOptions.tracks now only accepts 'all' | 'primary', no longer
accepts an array of InputTrack instances.
4. Unnecessary sync getter wrappers (commit 8241820)
Removed deprecated sync getters for: hasOnlyKeyPackets, bitrate,
averageBitrate, isRelativeToUnixEpoch on InputTrack.
============================================================
API ADDITIONS (for context, not release-note-worthy on their own)
============================================================
- PacketRetrievalOptions.skipLiveWait (fixes #342)
- DiscardedTrack.trackOptions field
- TargetRequest.mimeType
- TrackDisposition.primary
- Source/Target/Output extend EventEmitter
- CanvasSink methods now async generators (getCanvas, canvases, canvasesAtTimestamps)
- EncodedPacketSink.getFirstPacket/getPacket/getNextPacket now async
- Muxer writers deferred to start() instead of constructor
+12 -73
View File
@@ -23,7 +23,9 @@
chunked: true,
chunkSize: 2**20
});
const outputFormat = new Mediabunny.Mp4OutputFormat();
const outputFormat = new Mediabunny.HlsOutputFormat({
segmentFormat: new Mediabunny.MpegTsOutputFormat(),
});
const p = document.createElement('p');
p.textContent = 'Capturing...';
@@ -50,7 +52,7 @@
const output = new Mediabunny.Output({
format: outputFormat,
target
target: new Mediabunny.PathedTarget('master.m3u8', ({ path }) => new Mediabunny.BufferTarget()),
});
let input;
@@ -119,76 +121,13 @@
bitrate: 320000
},
*/
video: (track) => ({
discard: true,
//discard: true,
//discard: !tracks.includes(track),
//forceTranscode: true,
//forceTranscode: true,
//width: 1280,
//forceTranscode: true,
//allowRotationMetadata: false,
//width: 720,
//frameRate: 30,
//bitrate: Mediabunny.QUALITY_VERY_LOW,
//discard: true,
/*
process: (sample) => {
if (!ctx) {
// Create a canvas for image compositing
const canvas = new OffscreenCanvas(
sample.displayWidth,
sample.displayHeight,
);
ctx = canvas.getContext('2d');
}
console.log(ctx.canvas.width, ctx.canvas.height);
ctx.clearRect(0, 0, ctx.canvas.width, ctx.canvas.height);
sample.drawWithFit(ctx, { fit: 'fill' });
//ctx.drawImage(watermark, 32, 32);
return ctx.canvas;
},
*/
//width: 300,
//alpha: 'keep',
//width: 320,
//discard: true,
//discard: true,
//crop: {
// left: 0,
// top: 0,
// width: 500,
// height: 500,
//},
//rotate: 90,
//width: 200,
//height: 500,
//fit: 'contain',
//forceTranscode: true,
//codec: 'avc',
//fit: 'contain',
//frameRate: 27.123,
//width: 320,
//forceTranscode: true,
//codec: 'av1',
//discard: true,
//width: 1280,
//discard: true,
//width: 640
//forceTranscode: true,
//rotate: 90
//width: 720 ?? 2160,
//height: 1280 ?? 3840,
//fit: 'contain',
//rotate: 90,
//width: 512,
//height: 512,
//width: 200,
//height: 100,
}),
video: [
{ height: 1080 },
{ height: 720 },
{ height: 480 },
{ height: 360 },
{ height: 240 },
],
tags: {} ?? {
title: 'Bigggy',
artist: 'Buck Bunny',
@@ -207,7 +146,7 @@
}
},
trim: {
end: 5,
//end: 5,
//start,
//end: start + 5,
////start: 0,
+23
View File
@@ -279,6 +279,29 @@ output.addAudioTrack(audioSourceDe, { languageCode: 'de', name: 'German' });
output.addAudioTrack(audioSourceEs, { languageCode: 'es', name: 'Spanish' });
```
### Track disposition
You can control how a track is selected by players using the `disposition` option. The following fields map to HLS attributes on `#EXT-X-MEDIA`:
- `disposition.primary` → `DEFAULT=YES` (only one per media rendition group is allowed; subsequent `primary` tracks in the same group are ignored)
- `disposition.default` → `AUTOSELECT=YES` (also implied when `primary` is `true`)
- `disposition.forced` → `FORCED=YES`
For example:
```ts
output.addAudioTrack(audioSourceEn, {
languageCode: 'en',
name: 'English',
disposition: { primary: true },
});
output.addAudioTrack(audioSourceDe, {
languageCode: 'es',
name: 'Spanish',
disposition: { default: false },
});
```
When unspecified, `default` is `true` and `primary` and `forced` are `false`.
### I-frame only tracks
A track registered like this will be emitted via `#EXT-X-I-FRAME-STREAM-INF`:
+6 -6
View File
@@ -1,12 +1,12 @@
{
"name": "mediabunny",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "mediabunny",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"license": "MPL-2.0",
"workspaces": [
"packages/*"
@@ -12077,7 +12077,7 @@
},
"packages/aac-encoder": {
"name": "@mediabunny/aac-encoder",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"license": "MPL-2.0",
"devDependencies": {
"@types/emscripten": "^1.40.1"
@@ -12092,7 +12092,7 @@
},
"packages/ac3": {
"name": "@mediabunny/ac3",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"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.5",
"version": "1.42.0-beta.6",
"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.5",
"version": "1.42.0-beta.6",
"license": "MPL-2.0",
"devDependencies": {
"@types/emscripten": "^1.40.1"
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "mediabunny",
"author": "Vanilagy",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"description": "Pure TypeScript media toolkit for reading, writing, and converting media files, directly in the browser.",
"type": "module",
"workspaces": [
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "@mediabunny/aac-encoder",
"author": "Vanilagy",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"description": "AAC encoder extension for Mediabunny, based on FFmpeg.",
"main": "./dist/bundles/mediabunny-aac-encoder.mjs",
"module": "./dist/bundles/mediabunny-aac-encoder.mjs",
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "@mediabunny/ac3",
"author": "Vanilagy",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"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 -1
View File
@@ -1,7 +1,7 @@
{
"name": "@mediabunny/flac-encoder",
"author": "Vanilagy",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"description": "FLAC encoder extension for Mediabunny, based on libFLAC.",
"main": "./dist/bundles/mediabunny-flac-encoder.mjs",
"module": "./dist/bundles/mediabunny-flac-encoder.mjs",
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "@mediabunny/mp3-encoder",
"author": "Vanilagy",
"version": "1.42.0-beta.5",
"version": "1.42.0-beta.6",
"description": "MP3 encoder extension for Mediabunny, based on LAME.",
"main": "./dist/bundles/mediabunny-mp3-encoder.mjs",
"module": "./dist/bundles/mediabunny-mp3-encoder.mjs",
+12 -3
View File
@@ -898,6 +898,11 @@ export class HlsMuxer extends Muxer {
target: new PathedTarget(
fullSegmentPath,
async (request: TargetRequest) => {
const proxiedRequest: TargetRequest = {
...request,
isRoot: false,
};
if (request.isRoot) {
if (playlist.singleFile) {
const slice = playlist.singleFile.target.slice(playlist.singleFile.nextOffset);
@@ -905,7 +910,7 @@ export class HlsMuxer extends Muxer {
return slice;
} else {
const target = await this.output._getTarget(request);
const target = await this.output._getTarget(proxiedRequest);
outputTarget = target;
target.on('write', ({ end }) => segmentSize = Math.max(segmentSize, end));
@@ -913,7 +918,7 @@ export class HlsMuxer extends Muxer {
}
}
return this.output._getTarget(request);
return this.output._getTarget(proxiedRequest);
},
),
initTarget: async () => {
@@ -955,8 +960,12 @@ export class HlsMuxer extends Muxer {
info: null,
};
const fullInitPath = joinPaths(
joinPaths(pathedTarget.rootPath, playlist.path),
initPath,
);
const target = await this.output._getTarget({
path: initPath,
path: fullInitPath,
isRoot: false,
mimeType: playlist.segmentFormat.mimeType,
});
+8 -3
View File
@@ -11,7 +11,7 @@ import { ENCRYPTION_KEY_CACHE_GROUP, Input } from '../input';
import { Segment, SegmentedInput, SegmentedInputTrackDeclaration, SegmentRetrievalOptions } from '../segmented-input';
import { toDataView, joinPaths, last, assert, binarySearchLessOrEqual, arrayArgmin, wait } from '../misc';
import { readAllLines, readBytes, Reader } from '../reader';
import { CustomPathedSource, ReadableStreamSource, SourceRef } from '../source';
import { CustomPathedSource, ReadableStreamSource, SourceRef, SourceRequest } from '../source';
import { HlsDemuxer } from './hls-demuxer';
import {
AttributeList,
@@ -568,6 +568,11 @@ export class HlsSegmentedInput extends SegmentedInput {
async (request) => {
assert(request.isRoot); // Shouldn't fail since we don't allow recursive HLS
const proxiedRequest: SourceRequest = {
...request,
isRoot: false,
};
let ref: SourceRef;
const needsSlice = hlsSegment.location.offset > 0 || hlsSegment.location.length !== null;
@@ -576,7 +581,7 @@ export class HlsSegmentedInput extends SegmentedInput {
|| hlsSegment.encryption.method === 'SAMPLE-AES'
|| hlsSegment.encryption.method === 'SAMPLE-AES-CTR'
) {
ref = await this.input._getSourceCached(request);
ref = await this.input._getSourceCached(proxiedRequest);
if (needsSlice) {
const slice = ref.source.slice(
@@ -591,7 +596,7 @@ export class HlsSegmentedInput extends SegmentedInput {
const encryption = hlsSegment.encryption;
assert(encryption.iv);
let ciphertextRef = await this.input._getSourceCached(request);
let ciphertextRef = await this.input._getSourceCached(proxiedRequest);
if (needsSlice) {
// Slice before decrypting
const slice = ciphertextRef.source.slice(
+2 -1
View File
@@ -129,16 +129,17 @@ export {
Target,
TargetEvents,
TargetRequest,
AppendOnlyStreamTarget,
BufferTarget,
BufferTargetOptions,
FilePathTarget,
FilePathTargetOptions,
NullTarget,
PathedTarget,
RangedTarget,
StreamTarget,
StreamTargetOptions,
StreamTargetChunk,
PathedTarget,
} from './target';
export {
AnyIterable,
+14 -18
View File
@@ -378,7 +378,7 @@ export class Output<
/** @internal */
_muxer: Muxer;
/** @internal */
_targets = new Set<Target>();
_unfinalizedTargets = new Set<Target>();
/** @internal */
_rootWriterPromise: Promise<Writer> | null = null;
/** @internal */
@@ -439,12 +439,7 @@ export class Output<
throw new TypeError('options.target must be a Target or a PathedTarget.');
}
if (options.target instanceof Target) {
if (options.target._output) {
throw new Error('Target is already used for another output.');
}
options.target._output = this;
this._targets.add(options.target);
this._rememberTarget(options.target);
}
if (
options.initTarget !== undefined
@@ -466,8 +461,7 @@ export class Output<
this._initTarget = options.initTarget ?? null;
if (this._initTarget instanceof Target) {
this._initTarget._output = this;
this._targets.add(this._initTarget);
this._rememberTarget(this._initTarget);
}
this._muxer = options.format._createMuxer(this);
@@ -498,18 +492,23 @@ export class Output<
assert(this._target instanceof PathedTarget);
const target = await this._getTargetValidated(request);
target._output = this;
this._emit('target', { target, request, isRoot: request.isRoot });
if (this.state === 'canceled') {
await target._close();
} else {
this._targets.add(target);
this._rememberTarget(target);
}
return target;
}
/** @internal */
_rememberTarget(target: Target) {
this._unfinalizedTargets.add(target);
target.on('finalized', () => this._unfinalizedTargets.delete(target), { once: true });
}
/** @internal */
async _getInitTarget(): Promise<T> {
assert(this._initTarget !== null);
@@ -519,12 +518,11 @@ export class Output<
}
const target = await this._initTarget();
target._output = this;
if (this.state === 'canceled') {
await target._close();
} else {
this._targets.add(target);
this._rememberTarget(target);
}
return target;
@@ -558,13 +556,11 @@ export class Output<
const result = this._getTargetValidated(request);
const handleResult = (target: T) => {
target._output = this;
if (this.state === 'canceled') {
// Promise thrown away here, but no way to surface it to the user really
void target._close();
} else {
this._targets.add(target);
this._rememberTarget(target);
}
this._emit('target', { target, request, isRoot: true });
@@ -849,8 +845,8 @@ export class Output<
const promises = this._tracks.map(x => x.source._flushOrWaitForOngoingClose(true)); // Force close
await Promise.all(promises);
await Promise.all([...this._targets].map(target => target._close()));
this._targets.clear();
await Promise.all([...this._unfinalizedTargets].map(target => target._close()));
this._unfinalizedTargets.clear();
} finally {
release();
}
+31 -6
View File
@@ -558,7 +558,7 @@ export class BlobSource extends Source {
}
}
worker.running = false;
this._orchestrator.signalWorkerStoppedRunning(worker);
if (worker.aborted) {
// MDN: "Calling this method signals a loss of interest in the stream by a consumer."
@@ -853,7 +853,7 @@ export class UrlSource extends PathedSource {
while (true) {
if (worker.currentPos >= worker.targetPos || worker.aborted) {
abortController.abort();
worker.running = false;
this._orchestrator.signalWorkerStoppedRunning(worker);
return;
}
@@ -1201,7 +1201,7 @@ export class StreamSource extends Source {
}
}
worker.running = false;
this._orchestrator.signalWorkerStoppedRunning(worker);
}
/** @internal */
@@ -1608,6 +1608,8 @@ class ReadOrchestrator {
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null> {
assert(!this.disposed);
const prefetchRange = this.options.prefetchProfile(innerStart, innerEnd, this.workers);
const outerStart = Math.max(prefetchRange.start, minReadPosition);
const outerEnd = Math.min(prefetchRange.end, this.fileSize ?? Infinity, maxReadPosition);
@@ -1868,7 +1870,11 @@ class ReadOrchestrator {
for (let i = 0; i < this.workers.length; i++) {
const worker = this.workers[i]!;
if (!worker.running && (!oldestWorker || worker.age < oldestWorker.age)) {
if (
!worker.running
&& worker.pendingSlices.length === 0
&& (!oldestWorker || worker.age < oldestWorker.age)
) {
oldestIndex = i;
oldestWorker = worker;
}
@@ -1952,13 +1958,18 @@ class ReadOrchestrator {
// Here we merge everything into one "megaworker" that spans the entire file. We assume the passed-in worker
// is already configured to be a megaworker.
const uniqueSlices = new Set(worker.pendingSlices);
for (let i = 0; i < this.workers.length; i++) {
const otherWorker = this.workers[i]!;
if (otherWorker === worker) {
continue;
}
worker.pendingSlices.push(...otherWorker.pendingSlices);
for (const slice of otherWorker.pendingSlices) {
uniqueSlices.add(slice);
}
otherWorker.aborted = true;
otherWorker.pendingSlices.length = 0;
this.workers.splice(i, 1);
@@ -1967,9 +1978,13 @@ class ReadOrchestrator {
for (let i = 0; i < this.queuedReads.length; i++) {
const queuedRead = this.queuedReads[i]!;
worker.pendingSlices.push(...queuedRead.pendingSlices);
for (const slice of queuedRead.pendingSlices) {
uniqueSlices.add(slice);
}
}
worker.pendingSlices = [...uniqueSlices];
this.queuedReads.length = 0;
}
@@ -2101,6 +2116,16 @@ class ReadOrchestrator {
}
}
signalWorkerStoppedRunning(worker: ReadWorker) {
worker.running = false;
// When a worker stops running, that means it has hit its targetPos. It might still have pendingSlices assigned,
// but this is because those pending slices cover data that other workers are assigned to fill. Since targetPos
// has been reached, we can confidently say that this worker has completed its share of work on the pending
// slices and must no longer care about them.
worker.pendingSlices.length = 0;
}
/** Called when a worker reaches the end of the underlying data and must be cleaned up. */
onWorkerFinished(worker: ReadWorker) {
const index = this.workers.indexOf(worker);
+12 -3
View File
@@ -7,7 +7,6 @@
*/
import type { FileHandle } from 'node:fs/promises';
import { Output } from './output';
import * as nodeAlias from './node';
import { assert, EventEmitter, FilePath, MaybePromise } from './misc';
@@ -39,7 +38,7 @@ export type TargetEvents = {
*/
export abstract class Target extends EventEmitter<TargetEvents> {
/** @internal */
_output: Output | null = null;
_writerAcquired = false;
/** @internal */
_monotonicity: boolean | null = null; // null = unknown
@@ -582,6 +581,17 @@ export class StreamTarget extends Target {
}
}
/**
* This target writes to a `WritableStream<Uint8Array>`, meaning all writes are necessarily append-only and involve no
* seeking. Great for streaming data to a source that can only accept sequential data, like an HTTP server processing
* an incoming upload.
*
* Note that using this target *requires* that the underlying format write data sequentially. Not all formats do this,
* and this target will throw for the formats that don't. Check the guide for more.
*
* @group Output targets
* @public
*/
export class AppendOnlyStreamTarget extends Target {
/** @internal */
_writable: WritableStream<Uint8Array>;
@@ -787,7 +797,6 @@ export class RangedTarget extends Target {
this._baseTarget = baseTarget;
this._offset = offset;
this._output = baseTarget._output;
}
/** @internal */
+5 -2
View File
@@ -17,8 +17,13 @@ export class Writer {
private pos = 0;
constructor(target: Target, isMonotonic: boolean) {
if (target._writerAcquired) {
throw new Error('Can\'t have multiple Writers for the same Target.');
}
this.target = target;
target._setMonotonicity(isMonotonic);
target._writerAcquired = true;
}
start() {
@@ -55,10 +60,8 @@ export class Writer {
/** Called after muxing has finished. */
async finalize() {
assert(this.started && !this.finalized);
assert(this.target._output);
await this.target._finalize();
this.target._output._targets.delete(this.target);
this.finalized = true;
}
+21
View File
@@ -948,3 +948,24 @@ test.concurrent('Widevine encryption (SAMPLE-AES-CTR) succeeds with buffer keys'
assert(lastPacket);
expect(lastPacket.timestamp + lastPacket.duration).toBe(60);
});
test.concurrent('SourceRequest.isRoot', async () => {
using input = new Input({
source: new CustomPathedSource(
'https://test-streams.mux.dev/x36xhzz/x36xhzz.m3u8',
({ path, isRoot }) => {
if (isRoot) {
expect(path).toBe('https://test-streams.mux.dev/x36xhzz/x36xhzz.m3u8');
}
return new UrlSource(path);
},
),
formats: ALL_FORMATS,
});
const videoTrack = await input.getPrimaryVideoTrack();
assert(videoTrack);
await videoTrack.computeDuration();
});
+34
View File
@@ -2903,3 +2903,37 @@ test('Append-only stream with monotonicity violation', async () => {
await expect(source.add(new EncodedPacket(avcPacketData, 'key', 2, 0), avcMetadata))
.rejects.toThrow('AppendOnlyStreamTarget');
});
test('Relative paths & isRoot', async () => {
const output = new Output({
format: new HlsOutputFormat({
segmentFormat: new CmafOutputFormat(),
getPlaylistPath: () => `a/folder/playlist.m3u8`,
}),
target: new PathedTarget('path/to/master.m3u8', (request) => {
if (request.isRoot) {
expect(request.path).toBe('path/to/master.m3u8');
} else {
expect(request.path.startsWith('path/to/a/folder/')).toBe(true);
}
return new BufferTarget();
}),
});
const source = videoSource();
output.addVideoTrack(source);
await output.start();
await source.add(new EncodedPacket(avcPacketData, 'key', 0, 0), avcMetadata);
await source.add(new EncodedPacket(avcPacketData, 'delta', 0.5, 0), avcMetadata);
await source.add(new EncodedPacket(avcPacketData, 'delta', 1, 0), avcMetadata);
await source.add(new EncodedPacket(avcPacketData, 'delta', 1.5, 0), avcMetadata);
await source.add(new EncodedPacket(avcPacketData, 'key', 2, 0), avcMetadata);
await source.add(new EncodedPacket(avcPacketData, 'delta', 2.5, 0), avcMetadata);
await source.add(new EncodedPacket(avcPacketData, 'delta', 3, 0), avcMetadata);
await source.add(new EncodedPacket(avcPacketData, 'delta', 3.5, 0), avcMetadata);
await output.finalize();
});