Merge remote-tracking branch 'origin/main' into mpeg-ts

This commit is contained in:
Vanilagy
2026-01-16 21:04:14 +01:00
21 changed files with 297 additions and 86 deletions
+21 -2
View File
@@ -822,7 +822,7 @@ export class Conversion {
}
if (this._canceled) {
await new Promise(() => {}); // Never resolve
throw new ConversionCanceledError();
}
await this.output.finalize();
@@ -832,7 +832,10 @@ export class Conversion {
}
}
/** Cancels the conversion process. Does nothing if the conversion is already complete. */
/**
* Cancels the conversion process, causing any ongoing `execute` call to throw a `ConversionCanceledError`.
* Does nothing if the conversion is already complete.
*/
async cancel() {
if (this.output.state === 'finalizing' || this.output.state === 'finalized') {
return;
@@ -1159,6 +1162,7 @@ export class Conversion {
for await (const sample of sink.samples(this._startTimestamp, this._endTimestamp)) {
if (this._canceled) {
sample.close();
lastSample?.close();
return;
}
@@ -1441,6 +1445,7 @@ export class Conversion {
const sink = new AudioSampleSink(track);
for await (const sample of sink.samples(undefined, this._endTimestamp)) {
if (this._canceled) {
sample.close();
return;
}
@@ -1551,6 +1556,7 @@ export class Conversion {
for await (const sample of iterator) {
if (this._canceled) {
sample.close();
return;
}
@@ -1589,6 +1595,19 @@ export class Conversion {
}
}
/**
* Thrown when a conversion couldn't complete due to being canceled.
* @group Conversion
* @public
*/
export class ConversionCanceledError extends Error {
/** Creates a new {@link ConversionCanceledError}. */
constructor(message = 'Conversion has been canceled.') {
super(message);
this.name = 'ConversionCanceledError';
}
}
const MAX_TIMESTAMP_GAP = 5;
/**
+1
View File
@@ -194,6 +194,7 @@ export {
ConversionOptions,
ConversionVideoOptions,
ConversionAudioOptions,
ConversionCanceledError,
DiscardedTrack,
} from './conversion';
export {
+2 -2
View File
@@ -156,7 +156,7 @@ export class MatroskaInputFormat extends InputFormat {
}
const dataSize = readElementSize(headerSlice);
if (dataSize === null) {
if (typeof dataSize !== 'number') {
return false; // Miss me with that shit
}
@@ -172,7 +172,7 @@ export class MatroskaInputFormat extends InputFormat {
const { id, size } = header;
const dataStartPos = dataSlice.filePos;
if (size === null) return false;
if (size === undefined) return false;
switch (id) {
case EBMLId.EBMLVersion: {
+55 -18
View File
@@ -7,7 +7,7 @@
*/
import { MediaCodec } from '../codec';
import { assertNever, textDecoder, textEncoder } from '../misc';
import { assert, assertNever, textDecoder, textEncoder } from '../misc';
import { FileSlice, readBytes, Reader, readF32Be, readF64Be, readU8 } from '../reader';
import { Writer } from '../writer';
@@ -470,6 +470,10 @@ export const MIN_HEADER_SIZE = 2; // 1-byte ID and 1-byte size
export const MAX_HEADER_SIZE = 2 * MAX_VAR_INT_SIZE; // 8-byte ID and 8-byte size
export const readVarIntSize = (slice: FileSlice) => {
if (slice.remainingLength < 1) {
return null;
}
const firstByte = readU8(slice);
slice.skip(-1);
@@ -484,10 +488,19 @@ export const readVarIntSize = (slice: FileSlice) => {
mask >>= 1;
}
// Check if we have enough bytes to read the full varint
if (slice.remainingLength < width) {
return null;
}
return width;
};
export const readVarInt = (slice: FileSlice) => {
if (slice.remainingLength < 1) {
return null;
}
// Read the first byte to determine the width of the variable-length integer
const firstByte = readU8(slice);
@@ -503,6 +516,11 @@ export const readVarInt = (slice: FileSlice) => {
mask >>= 1;
}
if (slice.remainingLength < width - 1) {
// Not enough bytes
return null;
}
// First byte's value needs the marker bit cleared
let value = firstByte & (mask - 1);
@@ -563,39 +581,58 @@ export const readElementId = (slice: FileSlice) => {
return null;
}
if (slice.remainingLength < size) {
return null; // It don't fit
}
const id = readUnsignedInt(slice, size);
return id;
};
export const readElementSize = (slice: FileSlice) => {
let size: number | null = readU8(slice);
/** Returns `undefined` to indicate the EBML undefined size. Returns `null` if the size couldn't be read. */
export const readElementSize = (slice: FileSlice): number | undefined | null => {
// Need at least 1 byte to read the size
if (slice.remainingLength < 1) {
return null;
}
if (size === 0xff) {
size = null;
} else {
slice.skip(-1);
size = readVarInt(slice);
const firstByte = readU8(slice);
// In some (livestreamed) files, this is the value of the size field. While this technically is just a very
// large number, it is intended to behave like the reserved size 0xFF, meaning the size is undefined. We
// catch the number here. Note that it cannot be perfectly represented as a double, but the comparison works
// nonetheless.
// eslint-disable-next-line no-loss-of-precision
if (size === 0x00ffffffffffffff) {
size = null;
}
if (firstByte === 0xff) {
return undefined;
}
slice.skip(-1);
const size = readVarInt(slice);
if (size === null) {
return null;
}
// In some (livestreamed) files, this is the value of the size field. While this technically is just a very
// large number, it is intended to behave like the reserved size 0xFF, meaning the size is undefined. We
// catch the number here. Note that it cannot be perfectly represented as a double, but the comparison works
// nonetheless.
// eslint-disable-next-line no-loss-of-precision
if (size === 0x00ffffffffffffff) {
return undefined;
}
return size;
};
export const readElementHeader = (slice: FileSlice) => {
assert(slice.remainingLength >= MIN_HEADER_SIZE);
const id = readElementId(slice);
if (id === null) {
return null;
}
const size = readElementSize(slice);
if (size === null) {
return null;
}
return { id, size };
};
@@ -720,8 +757,8 @@ export const CODEC_STRING_MAP: Partial<Record<MediaCodec, string>> = {
'webvtt': 'S_TEXT/WEBVTT',
};
export function assertDefinedSize(size: number | null): asserts size is number {
if (size === null) {
export function assertDefinedSize(size: number | undefined): asserts size is number {
if (size === undefined) {
throw new Error('Undefined element size is used in a place where it is not supported.');
}
};
+23 -14
View File
@@ -193,6 +193,7 @@ type InternalTrack = {
codecId: string | null;
codecPrivate: Uint8Array | null;
defaultDuration: number | null;
defaultDurationNs: number | null;
name: string | null;
languageCode: string;
decodingInstructions: DecodingInstruction[];
@@ -346,7 +347,7 @@ export class MatroskaDemuxer extends Demuxer {
} else if (id === EBMLId.Segment) { // Segment found!
await this.readSegment(dataStartPos, size);
if (size === null) {
if (size === undefined) {
// Segment sizes can be undefined (common in livestreamed files), so assume this is the last
// and only segment
break;
@@ -364,7 +365,7 @@ export class MatroskaDemuxer extends Demuxer {
// doesn't contain any of the clusters that follow it. In the case, we apply the following logic: if
// we find a top-level cluster, attribute it to the previous segment.
if (size === null) {
if (size === undefined) {
// Just in case this is one of those weird sizeless clusters, let's do our best and still try to
// determine its size.
const nextElementPos = await searchForNextElementId(
@@ -389,7 +390,7 @@ export class MatroskaDemuxer extends Demuxer {
})();
}
async readSegment(segmentDataStart: number, dataSize: number | null) {
async readSegment(segmentDataStart: number, dataSize: number | undefined) {
this.currentSegment = {
seekHeadSeen: false,
infoSeen: false,
@@ -406,7 +407,7 @@ export class MatroskaDemuxer extends Demuxer {
cuePoints: [],
dataStartPos: segmentDataStart,
elementEndPos: dataSize === null
elementEndPos: dataSize === undefined
? null // Assume it goes until the end of the file
: segmentDataStart + dataSize,
clusterSeekStartPos: segmentDataStart,
@@ -483,7 +484,7 @@ export class MatroskaDemuxer extends Demuxer {
break; // Stop at the first cluster
}
if (size === null) {
if (size === undefined) {
break;
} else {
currentPos = dataStartPos + size;
@@ -536,6 +537,13 @@ export class MatroskaDemuxer extends Demuxer {
this.currentSegment.timestampFactor = 1e9 / 1e6;
}
// Compute default duration for all tracks now that we have the timestamp factor
for (const track of this.currentSegment.tracks) {
if (track.defaultDurationNs !== null) {
track.defaultDuration = (this.currentSegment.timestampFactor * track.defaultDurationNs) / 1e9;
}
}
// Put default tracks first
this.currentSegment.tracks.sort((a, b) => Number(b.disposition.default) - Number(a.disposition.default));
@@ -606,7 +614,7 @@ export class MatroskaDemuxer extends Demuxer {
let size = elementHeader.size;
const dataStartPos = headerSlice.filePos;
if (size === null) {
if (size === undefined) {
// The cluster's size is undefined (can happen in livestreamed files). We'd still like to know the size of
// it, so we have no other choice but to iterate over the EBML structure until we find an element at level
// 0 or 1, indicating the end of the cluster (all elements inside the cluster are at level 2).
@@ -908,9 +916,7 @@ export class MatroskaDemuxer extends Demuxer {
}
readContiguousElements(slice: FileSlice, stopIds?: number[]) {
const startIndex = slice.filePos;
while (slice.filePos - startIndex <= slice.length - MIN_HEADER_SIZE) {
while (slice.remainingLength >= MIN_HEADER_SIZE) {
const startPos = slice.filePos;
const foundElement = this.traverseElement(slice, stopIds);
@@ -996,6 +1002,7 @@ export class MatroskaDemuxer extends Demuxer {
codecId: null,
codecPrivate: null,
defaultDuration: null,
defaultDurationNs: null,
name: null,
languageCode: UNDETERMINED_LANGUAGE,
decodingInstructions: [],
@@ -1005,6 +1012,11 @@ export class MatroskaDemuxer extends Demuxer {
this.readContiguousElements(slice.slice(dataStartPos, size));
// Check if track was disabled during parsing (e.g., by FlagEnabled being 0)
if (!this.currentTrack) {
break;
}
if (this.currentTrack.decodingInstructions.some((instruction) => {
return instruction.data?.type !== 'decompress'
|| instruction.scope !== ContentEncodingScope.Block
@@ -1149,7 +1161,6 @@ export class MatroskaDemuxer extends Demuxer {
const enabled = readUnsignedInt(slice, size);
if (!enabled) {
this.currentSegment!.tracks.pop();
this.currentTrack = null;
}
}; break;
@@ -1204,9 +1215,7 @@ export class MatroskaDemuxer extends Demuxer {
case EBMLId.DefaultDuration: {
if (!this.currentTrack) break;
this.currentTrack.defaultDuration
= this.currentTrack.segment.timestampFactor * readUnsignedInt(slice, size) / 1e9;
this.currentTrack.defaultDurationNs = readUnsignedInt(slice, size);
}; break;
case EBMLId.Name: {
@@ -2223,7 +2232,7 @@ abstract class MatroskaTrackBacking implements InputTrackBacking {
}
}
if (size === null) {
if (size === undefined) {
// Undefined element size (can happen in livestreamed files). In this case, we need to do some
// searching to determine the actual size of the element.
+7 -17
View File
@@ -480,23 +480,13 @@ export abstract class BaseMediaSampleSink<
let currentPacket: EncodedPacket | null = keyPacket;
let endPacket: EncodedPacket | undefined = undefined;
if (endTimestamp < Infinity) {
// When an end timestamp is set, we cannot simply use that for the packet iterator due to out-of-order
// frames (B-frames). Instead, we'll need to keep decoding packets until we get a frame that exceeds
// this end time. However, we can still put a bound on it: Since key frames are by definition never
// out of order, we can stop at the first key frame after the end timestamp.
const packet = await packetSink.getPacket(endTimestamp);
const keyPacket = !packet
? null
: packet.type === 'key' && packet.timestamp === endTimestamp
? packet
: await packetSink.getNextKeyPacket(packet, { verifyKeyPackets: true });
if (keyPacket) {
endPacket = keyPacket;
}
}
// B-frames make it exceedingly difficult to properly define an upper bound for packet iteration if an end
// timestamp is set, so we just don't do it. The case that makes it especially tricky is when the frames
// following a key frame have a lower timestamp than the keyframe; something that quite frequently happens
// in HEVC streams. The price to pay for not upper-bounding the packet iterator is a slight increase in
// decoder work at the end of the range, but the added correctness and reliability makes this tradeoff worth
// it.
const endPacket = undefined;
const packets = packetSink.packets(keyPacket ?? undefined, endPacket);
await packets.next(); // Skip the start packet as we already have it
+4 -6
View File
@@ -134,12 +134,10 @@ export abstract class MediaSource {
/** @internal */
async _flushOrWaitForOngoingClose(forceClose: boolean) {
if (this._closingPromise) {
// Since closing also flushes, we don't want to do it twice
return this._closingPromise;
} else {
return this._flushAndClose(forceClose);
}
return this._closingPromise ??= (async () => {
await this._flushAndClose(forceClose);
this._closed = true;
})();
}
}
+4
View File
@@ -322,6 +322,10 @@ export class OggDemuxer extends Demuxer {
}
const totalPacketSize = chunks.reduce((sum, chunk) => sum + chunk.length, 0);
if (totalPacketSize === 0) {
return null; // Invalid packet, treat it as end of stream
}
const packetData = new Uint8Array(totalPacketSize);
let offset = 0;
+34 -17
View File
@@ -47,12 +47,13 @@ type OggTrackData = {
currentPageData: Uint8Array[];
currentPageSize: number;
currentPageStartsWithFreshPacket: boolean;
currentPageStartTimestampInSamples: number;
};
type Packet = {
data: Uint8Array;
endGranulePosition: number;
timestamp: number;
timestampInSamples: number;
durationInSamples: number;
forcePageFlush: boolean;
};
@@ -132,6 +133,7 @@ export class OggMuxer extends Muxer {
currentPageData: [],
currentPageSize: 27,
currentPageStartsWithFreshPacket: true,
currentPageStartTimestampInSamples: 0,
};
this.queueHeaderPackets(newTrackData, meta);
@@ -199,18 +201,18 @@ export class OggMuxer extends Muxer {
trackData.packetQueue.push({
data: identificationHeader,
endGranulePosition: 0,
timestamp: 0,
timestampInSamples: 0,
durationInSamples: 0,
forcePageFlush: true,
}, {
data: commentHeader,
endGranulePosition: 0,
timestamp: 0,
timestampInSamples: 0,
durationInSamples: 0,
forcePageFlush: false,
}, {
data: setupHeader,
endGranulePosition: 0,
timestamp: 0,
timestampInSamples: 0,
durationInSamples: 0,
forcePageFlush: true, // The last header packet must flush the page
});
@@ -239,13 +241,13 @@ export class OggMuxer extends Muxer {
trackData.packetQueue.push({
data: identificationHeader,
endGranulePosition: 0,
timestamp: 0,
timestampInSamples: 0,
durationInSamples: 0,
forcePageFlush: true,
}, {
data: commentHeader,
endGranulePosition: 0,
timestamp: 0,
timestampInSamples: 0,
durationInSamples: 0,
forcePageFlush: true, // The last header packet must flush the page
});
@@ -275,8 +277,8 @@ export class OggMuxer extends Muxer {
trackData.packetQueue.push({
data: packet.data,
endGranulePosition: trackData.currentTimestampInSamples,
timestamp: currentTimestampInSamples / trackData.internalSampleRate,
timestampInSamples: currentTimestampInSamples,
durationInSamples,
forcePageFlush: false,
});
@@ -338,10 +340,10 @@ export class OggMuxer extends Muxer {
if (
trackData.packetQueue.length > 0
&& trackData.packetQueue[0]!.timestamp < minTimestamp
&& trackData.packetQueue[0]!.timestampInSamples < minTimestamp
) {
trackWithMinTimestamp = trackData;
minTimestamp = trackData.packetQueue[0]!.timestamp;
minTimestamp = trackData.packetQueue[0]!.timestampInSamples;
}
}
@@ -361,6 +363,20 @@ export class OggMuxer extends Muxer {
}
writePacket(trackData: OggTrackData, packet: Packet, isFinalPacket: boolean) {
const packetEndTimestampInSamples = packet.timestampInSamples + packet.durationInSamples;
if (this.format._options.maximumPageDuration !== undefined) {
const maxDurationInSamples = this.format._options.maximumPageDuration * trackData.internalSampleRate;
if (
trackData.currentLacingValues.length > 0
&& packetEndTimestampInSamples - trackData.currentPageStartTimestampInSamples > maxDurationInSamples
) {
// Flush the current page early to avoid exceeding the maximum page duration
this.writePage(trackData, false);
}
}
let remainingLength = packet.data.length;
let dataStartOffset = 0;
let dataOffset = 0;
@@ -401,7 +417,7 @@ export class OggMuxer extends Muxer {
const slice = packet.data.subarray(dataStartOffset);
trackData.currentPageData.push(slice);
trackData.currentPageSize += slice.length;
trackData.currentGranulePosition = packet.endGranulePosition;
trackData.currentGranulePosition = packetEndTimestampInSamples;
if (trackData.currentPageSize >= PAGE_SIZE_TARGET || packet.forcePageFlush) {
this.writePage(trackData, isFinalPacket);
@@ -452,6 +468,7 @@ export class OggMuxer extends Muxer {
trackData.currentPageData.length = 0;
trackData.currentPageSize = 27;
trackData.currentPageStartsWithFreshPacket = true;
trackData.currentPageStartTimestampInSamples = trackData.currentGranulePosition;
if (this.format._options.onPage) {
this.writer.startTrackingWrites();
+12
View File
@@ -730,6 +730,12 @@ export class WavOutputFormat extends OutputFormat {
* @public
*/
export type OggOutputFormatOptions = {
/**
* The maximum duration of each Ogg page, in seconds. This is useful for streaming contexts where more frequent page
* output is desired. By default, pages are only flushed when they exceed a certain size.
*/
maximumPageDuration?: number;
/**
* Will be called for each Ogg page that is written.
*
@@ -754,6 +760,12 @@ export class OggOutputFormat extends OutputFormat {
if (!options || typeof options !== 'object') {
throw new TypeError('options must be an object.');
}
if (
options.maximumPageDuration !== undefined
&& (!Number.isFinite(options.maximumPageDuration) || options.maximumPageDuration <= 0)
) {
throw new TypeError('options.maximumPageDuration, when provided, must be a positive number.');
}
if (options.onPage !== undefined && typeof options.onPage !== 'function') {
throw new TypeError('options.onPage, when provided, must be a function.');
}