From 95cc8beb06550a17d11940163b39025f0fa98da9 Mon Sep 17 00:00:00 2001 From: Vanilagy <1696106+Vanilagy@users.noreply.github.com> Date: Wed, 26 Mar 2025 21:35:19 +0100 Subject: [PATCH] Fix chunked writing issues & drop ChunkedStreamTargetWriter --- dev/convert.html | 12 +++- src/target.ts | 4 +- src/writer.ts | 174 ++++++++++++++++++----------------------------- 3 files changed, 78 insertions(+), 112 deletions(-) diff --git a/dev/convert.html b/dev/convert.html index ec062dd..9518045 100644 --- a/dev/convert.html +++ b/dev/convert.html @@ -15,9 +15,15 @@ const file = fileInput.files[0]; const source = new Metamuxer.BlobSource(file); - const target = new Metamuxer.BufferTarget(); + const target = new Metamuxer.StreamTarget(new WritableStream({ + write: console.log + }), { + chunked: true, + chunkSize: 2**20 + }); const outputFormat = new Metamuxer.OggOutputFormat({ - onPage: console.log + //streamable: true + //fastStart: 'fragmented' }); const button = document.createElement('button'); @@ -75,7 +81,7 @@ }, trim: { start: 0, - end: 60 + //end: 60 }, computeProgress: true }); diff --git a/src/target.ts b/src/target.ts index cd1c99c..90059b3 100644 --- a/src/target.ts +++ b/src/target.ts @@ -1,4 +1,4 @@ -import { BufferTargetWriter, ChunkedStreamTargetWriter, StreamTargetWriter, Writer } from './writer'; +import { BufferTargetWriter, StreamTargetWriter, Writer } from './writer'; import { Output } from './output'; /** @@ -93,6 +93,6 @@ export class StreamTarget extends Target { /** @internal */ _createWriter() { - return this._options.chunked ? new ChunkedStreamTargetWriter(this) : new StreamTargetWriter(this); + return new StreamTargetWriter(this); } } diff --git a/src/writer.ts b/src/writer.ts index 2320074..ae7fd1e 100644 --- a/src/writer.ts +++ b/src/writer.ts @@ -175,9 +175,25 @@ export class BufferTargetWriter extends Writer { } } +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 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. + * last flush. This is useful for streaming applications, like piping the output to disk. When using the chunked mode, + * data will first be accumulated in larger chunks, and then the entire chunk will be flushed out at once when ready. */ export class StreamTargetWriter extends Writer { private pos = 0; @@ -187,13 +203,26 @@ export class StreamTargetWriter extends Writer { start: number; }[] = []; + private lastWriteEnd = 0; private lastFlushEnd = 0; private writer: WritableStreamDefaultWriter | null = null; + // These variables regard chunked mode: + private chunked: boolean; + 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[] = []; + constructor(target: StreamTarget) { super(); this.target = target; + + this.chunked = target._options.chunked ?? false; + this.chunkSize = target._options.chunkSize ?? DEFAULT_CHUNK_SIZE; } override start() { @@ -201,6 +230,12 @@ export class StreamTargetWriter extends Writer { } write(data: Uint8Array) { + if (this.pos > this.lastWriteEnd) { + const paddingBytesNeeded = this.pos - this.lastWriteEnd; + this.pos = this.lastWriteEnd; + this.write(new Uint8Array(paddingBytesNeeded)); + } + this.maybeTrackWrites(data); this.sections.push({ @@ -208,6 +243,8 @@ export class StreamTargetWriter extends Writer { start: this.pos, }); this.pos += data.byteLength; + + this.lastWriteEnd = Math.max(this.lastWriteEnd, this.pos); } seek(newPos: number) { @@ -219,6 +256,14 @@ export class StreamTargetWriter extends Writer { } async flush() { + if (this.pos > this.lastWriteEnd) { + // There's a "void" between the last written byte and the next byte we're about to write. Let's pad that + // void with zeroes explicitly. + const paddingBytesNeeded = this.pos - this.lastWriteEnd; + this.pos = this.lastWriteEnd; + this.write(new Uint8Array(paddingBytesNeeded)); + } + assert(this.writer); if (this.sections.length === 0) return; @@ -268,91 +313,25 @@ export class StreamTargetWriter extends Writer { await this.writer.ready; // Allow the writer to apply backpressure } - void this.writer.write({ - type: 'write', - data: chunk.data, - position: chunk.start, - }); + if (this.chunked) { + // Let's first gather the data into bigger chunks before writing it + this.writeDataIntoChunks(chunk.data, chunk.start); + this.tryToFlushChunks(); + } else { + // Write out the data immediately + 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(); - } - - async close() { - 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 | 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.maybeTrackWrites(data); - - 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); @@ -382,10 +361,10 @@ export class ChunkedStreamTargetWriter extends Writer { for (let i = 0; i < this.chunks.length - 1; i++) { this.chunks[i]!.shouldFlush = true; } - this.queueChunksForFlush(); + this.tryToFlushChunks(); } - // If the data didn't fit in one chunk, recurse with the remaining datas + // If the data didn't fit in one chunk, recurse with the remaining data if (toWrite.byteLength < data.byteLength) { this.writeDataIntoChunks(data.subarray(toWrite.byteLength), position + toWrite.byteLength); } @@ -433,7 +412,7 @@ export class ChunkedStreamTargetWriter extends Writer { return this.chunks.indexOf(chunk); } - private queueChunksForFlush(force = false) { + private tryToFlushChunks(force = false) { assert(this.writer); for (let i = 0; i < this.chunks.length; i++) { @@ -441,42 +420,23 @@ export class ChunkedStreamTargetWriter extends Writer { 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({ + void this.writer.write({ 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); + finalize() { + if (this.chunked) { + this.tryToFlushChunks(true); } - this.flushedChunkQueue.length = 0; - } - - async finalize() { assert(this.writer); - - this.queueChunksForFlush(true); - await this.flush(); - return this.writer.close(); }