Files
mediabunny/src/writer.ts
T
2024-12-08 15:19:56 +01:00

379 lines
10 KiB
TypeScript

import { ArrayBufferTarget, StreamTarget, StreamTargetChunk } from './target';
import { assert } from './misc';
export abstract class Writer {
/** Setting this to true will cause the writer to ensure data is written in a strictly monotonic, streamable way. */
ensureMonotonicity = false;
start() {}
/** Writes the given data to the target, at the current position. */
abstract write(data: Uint8Array): void;
/** Sets the current position for future writes to a new one. */
abstract seek(newPos: number): void;
/** Returns the current position. */
abstract getPos(): number;
/** Signals to the writer that it may be time to flush. */
abstract flush(): Promise<void>;
/** Called after muxing has finished. */
abstract finalize(): Promise<void>;
}
/**
* Writes to an ArrayBufferTarget. Maintains a growable internal buffer during the muxing process, which will then be
* written to the ArrayBufferTarget once the muxing finishes.
*/
export class ArrayBufferTargetWriter extends Writer {
private pos = 0;
private target: ArrayBufferTarget;
private buffer = new ArrayBuffer(2 ** 16);
private bytes = new Uint8Array(this.buffer);
private maxPos = 0;
constructor(target: ArrayBufferTarget) {
super();
this.target = target;
}
private ensureSize(size: number) {
let newLength = this.buffer.byteLength;
while (newLength < size) newLength *= 2;
if (newLength === this.buffer.byteLength) return;
const newBuffer = new ArrayBuffer(newLength);
const newBytes = new Uint8Array(newBuffer);
newBytes.set(this.bytes, 0);
this.buffer = newBuffer;
this.bytes = newBytes;
}
write(data: Uint8Array) {
this.ensureSize(this.pos + data.byteLength);
this.bytes.set(data, this.pos);
this.pos += data.byteLength;
this.maxPos = Math.max(this.maxPos, this.pos);
}
seek(newPos: number) {
this.pos = newPos;
}
getPos() {
return this.pos;
}
async flush() {}
async finalize() {
this.ensureSize(this.pos);
this.target.buffer = this.buffer.slice(0, Math.max(this.maxPos, this.pos));
}
getSlice(start: number, end: number) {
return this.bytes.slice(start, end);
}
}
/**
* Writes to a StreamTarget every time it is flushed, sending out all of the new data written since the
* last flush. This is useful for streaming applications, like piping the output to disk.
*/
export class StreamTargetWriter extends Writer {
private pos = 0;
private target: StreamTarget;
private sections: {
data: Uint8Array;
start: number;
}[] = [];
private lastFlushEnd = 0;
private writer: WritableStreamDefaultWriter<StreamTargetChunk> | null = null;
constructor(target: StreamTarget) {
super();
this.target = target;
}
override start() {
this.writer = this.target._writable.getWriter();
}
write(data: Uint8Array) {
this.sections.push({
data: data.slice(),
start: this.pos,
});
this.pos += data.byteLength;
}
seek(newPos: number) {
this.pos = newPos;
}
getPos() {
return this.pos;
}
async flush() {
assert(this.writer);
if (this.sections.length === 0) return;
const chunks: {
start: number;
size: number;
data?: Uint8Array;
}[] = [];
const sorted = [...this.sections].sort((a, b) => a.start - b.start);
chunks.push({
start: sorted[0]!.start,
size: sorted[0]!.data.byteLength,
});
// Figure out how many contiguous chunks we have
for (let i = 1; i < sorted.length; i++) {
const lastChunk = chunks[chunks.length - 1]!;
const section = sorted[i]!;
if (section.start <= lastChunk.start + lastChunk.size) {
lastChunk.size = Math.max(lastChunk.size, section.start + section.data.byteLength - lastChunk.start);
} else {
chunks.push({
start: section.start,
size: section.data.byteLength,
});
}
}
for (const chunk of chunks) {
chunk.data = new Uint8Array(chunk.size);
// Make sure to write the data in the correct order for correct overwriting
for (const section of this.sections) {
// Check if the section is in the chunk
if (chunk.start <= section.start && section.start < chunk.start + chunk.size) {
chunk.data.set(section.data, section.start - chunk.start);
}
}
if (this.ensureMonotonicity && chunk.start !== this.lastFlushEnd) {
throw new Error('Internal error: Monotonicity violation.');
}
if (this.writer.desiredSize !== null && this.writer.desiredSize <= 0) {
await this.writer.ready; // Allow the writer to apply backpressure
}
void this.writer.write({
type: 'write',
data: chunk.data,
position: chunk.start,
});
this.lastFlushEnd = chunk.start + chunk.data.byteLength;
}
this.sections.length = 0;
}
finalize() {
assert(this.writer);
return this.writer.close();
}
}
const DEFAULT_CHUNK_SIZE = 2 ** 24;
const MAX_CHUNKS_AT_ONCE = 2;
interface Chunk {
start: number;
written: ChunkSection[];
data: Uint8Array;
shouldFlush: boolean;
}
interface ChunkSection {
start: number;
end: number;
}
/**
* Writes to a StreamTarget using a chunked approach: Data is first buffered in memory until it reaches a large enough
* size, which is when it is piped to the StreamTarget. This is helpful for reducing the total amount of writes.
*/
export class ChunkedStreamTargetWriter extends Writer {
private pos = 0;
private target: StreamTarget;
private chunkSize: number;
/**
* The data is divided up into fixed-size chunks, whose contents are first filled in RAM and then flushed out.
* A chunk is flushed if all of its contents have been written.
*/
private chunks: Chunk[] = [];
private lastFlushEnd = 0;
private writer: WritableStreamDefaultWriter<StreamTargetChunk> | null = null;
private flushedChunkQueue: StreamTargetChunk[] = [];
constructor(target: StreamTarget) {
super();
this.target = target;
this.chunkSize = target._options?.chunkSize ?? DEFAULT_CHUNK_SIZE;
if (!Number.isInteger(this.chunkSize) || this.chunkSize < 2 ** 10) {
throw new Error('Invalid StreamTarget options: chunkSize must be an integer not smaller than 1024.');
}
}
override start() {
this.writer = this.target._writable.getWriter();
}
write(data: Uint8Array) {
this.writeDataIntoChunks(data, this.pos);
this.queueChunksForFlush();
this.pos += data.byteLength;
}
seek(newPos: number) {
this.pos = newPos;
}
getPos() {
return this.pos;
}
private writeDataIntoChunks(data: Uint8Array, position: number) {
// First, find the chunk to write the data into, or create one if none exists
let chunkIndex = this.chunks.findIndex(x => x.start <= position && position < x.start + this.chunkSize);
if (chunkIndex === -1) chunkIndex = this.createChunk(position);
const chunk = this.chunks[chunkIndex]!;
// Figure out how much to write to the chunk, and then write to the chunk
const relativePosition = position - chunk.start;
const toWrite = data.subarray(0, Math.min(this.chunkSize - relativePosition, data.byteLength));
chunk.data.set(toWrite, relativePosition);
// Create a section describing the region of data that was just written to
const section: ChunkSection = {
start: relativePosition,
end: relativePosition + toWrite.byteLength,
};
this.insertSectionIntoChunk(chunk, section);
// Queue chunk for flushing to target if it has been fully written to
if (chunk.written[0]!.start === 0 && chunk.written[0]!.end === this.chunkSize) {
chunk.shouldFlush = true;
}
// Make sure we don't hold too many chunks in memory at once to keep memory usage down
if (this.chunks.length > MAX_CHUNKS_AT_ONCE) {
// Flush all but the last chunk
for (let i = 0; i < this.chunks.length - 1; i++) {
this.chunks[i]!.shouldFlush = true;
}
this.queueChunksForFlush();
}
// If the data didn't fit in one chunk, recurse with the remaining datas
if (toWrite.byteLength < data.byteLength) {
this.writeDataIntoChunks(data.subarray(toWrite.byteLength), position + toWrite.byteLength);
}
}
private insertSectionIntoChunk(chunk: Chunk, section: ChunkSection) {
let low = 0;
let high = chunk.written.length - 1;
let index = -1;
// Do a binary search to find the last section with a start not larger than `section`'s start
while (low <= high) {
const mid = Math.floor(low + (high - low + 1) / 2);
if (chunk.written[mid]!.start <= section.start) {
low = mid + 1;
index = mid;
} else {
high = mid - 1;
}
}
// Insert the new section
chunk.written.splice(index + 1, 0, section);
if (index === -1 || chunk.written[index]!.end < section.start) index++;
// Merge overlapping sections
while (index < chunk.written.length - 1 && chunk.written[index]!.end >= chunk.written[index + 1]!.start) {
chunk.written[index]!.end = Math.max(chunk.written[index]!.end, chunk.written[index + 1]!.end);
chunk.written.splice(index + 1, 1);
}
}
private createChunk(includesPosition: number) {
const start = Math.floor(includesPosition / this.chunkSize) * this.chunkSize;
const chunk: Chunk = {
start,
data: new Uint8Array(this.chunkSize),
written: [],
shouldFlush: false,
};
this.chunks.push(chunk);
this.chunks.sort((a, b) => a.start - b.start);
return this.chunks.indexOf(chunk);
}
private queueChunksForFlush(force = false) {
assert(this.writer);
for (let i = 0; i < this.chunks.length; i++) {
const chunk = this.chunks[i]!;
if (!chunk.shouldFlush && !force) continue;
for (const section of chunk.written) {
if (this.ensureMonotonicity && chunk.start + section.start !== this.lastFlushEnd) {
throw new Error('Internal error: Monotonicity violation.');
}
this.flushedChunkQueue.push({
type: 'write',
data: chunk.data.subarray(section.start, section.end),
position: chunk.start + section.start,
});
this.lastFlushEnd = chunk.start + section.end;
}
this.chunks.splice(i--, 1);
}
}
async flush() {
assert(this.writer);
if (this.flushedChunkQueue.length === 0) return;
for (const chunk of this.flushedChunkQueue) {
if (this.writer.desiredSize !== null && this.writer.desiredSize <= 0) {
await this.writer.ready; // Allow the writer to apply backpressure
}
void this.writer.write(chunk);
}
this.flushedChunkQueue.length = 0;
}
async finalize() {
assert(this.writer);
this.queueChunksForFlush(true);
await this.flush();
return this.writer.close();
}
}