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();
}