Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
66 changes: 65 additions & 1 deletion packages/ducjs/src/opfs-import.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,18 @@ export type OpfsChunkImporter = {
export type DucStreamImportOptions = {
compressedChunkSize?: number;
writeBufferSize?: number;
progressIntervalMs?: number;
onProgress?: (progress: DucStreamImportProgress) => void;
};

export type DucStreamImportProgress = {
phase: "streaming" | "complete";
compressedBytes: number;
uncompressedBytes: number;
elapsedMs: number;
readWaitMs: number;
processingMs: number;
opfsWriteMs: number;
};

class BufferedImporterWriter {
Expand All @@ -21,6 +33,7 @@ class BufferedImporterWriter {
constructor(
private readonly importer: OpfsChunkImporter,
bufferSize: number,
private readonly onWrite?: (byteLength: number, elapsedMs: number) => void,
) {
if (!Number.isSafeInteger(bufferSize) || bufferSize <= 0) {
throw new RangeError("writeBufferSize must be a positive safe integer");
Expand Down Expand Up @@ -53,9 +66,19 @@ class BufferedImporterWriter {
return this.totalBytes;
}

get pendingByteLength(): number {
return this.offset;
}

get writtenByteLength(): number {
return this.totalBytes;
}

private flush(): void {
if (this.offset === 0) return;
const startedAt = performance.now();
this.importer.writeChunk(this.buffer.subarray(0, this.offset));
this.onWrite?.(this.offset, performance.now() - startedAt);
this.totalBytes += this.offset;
this.offset = 0;
}
Expand All @@ -76,10 +99,41 @@ export async function importDucStreamToOpfs(
throw new RangeError("compressedChunkSize must be a positive safe integer");
}

const startedAt = performance.now();
const progressIntervalMs = options.progressIntervalMs ?? 500;
if (!Number.isFinite(progressIntervalMs) || progressIntervalMs <= 0) {
throw new RangeError("progressIntervalMs must be a positive number");
}
let lastProgressAt = startedAt;
let compressedBytes = 0;
let readWaitMs = 0;
let processingMs = 0;
let opfsWriteMs = 0;
const writer = new BufferedImporterWriter(
importer,
options.writeBufferSize ?? DEFAULT_WRITE_BUFFER_SIZE,
(_byteLength, elapsedMs) => {
opfsWriteMs += elapsedMs;
},
);

const emitProgress = (phase: DucStreamImportProgress["phase"]): void => {
if (!options.onProgress) return;
const now = performance.now();
if (phase === "streaming" && now - lastProgressAt < progressIntervalMs) {
return;
}
options.onProgress({
phase,
compressedBytes,
uncompressedBytes: writer.writtenByteLength + writer.pendingByteLength,
elapsedMs: now - startedAt,
readWaitMs,
processingMs,
opfsWriteMs,
});
lastProgressAt = now;
};
const decompressor = new Decompress((chunk) => writer.write(chunk));
const reader = stream.getReader();
const header = new Uint8Array(SQLITE_HEADER.byteLength);
Expand All @@ -105,10 +159,14 @@ export async function importDucStreamToOpfs(

try {
while (true) {
const readStartedAt = performance.now();
const { done, value } = await reader.read();
readWaitMs += performance.now() - readStartedAt;
if (done) break;
if (!value?.byteLength) continue;
compressedBytes += value.byteLength;

const processingStartedAt = performance.now();
let valueOffset = 0;
if (isRawSqlite === null) {
const headerBytes = Math.min(
Expand All @@ -128,16 +186,22 @@ export async function importDucStreamToOpfs(
if (isRawSqlite !== null && valueOffset < value.byteLength) {
writeDetected(value.subarray(valueOffset));
}
processingMs += performance.now() - processingStartedAt;
emitProgress("streaming");
}

const finalProcessingStartedAt = performance.now();
if (isRawSqlite === null) {
isRawSqlite = false;
writeCompressed(header.subarray(0, headerLength));
}
if (!isRawSqlite) {
decompressor.push(new Uint8Array(), true);
}
return writer.finish();
const writtenBytes = writer.finish();
processingMs += performance.now() - finalProcessingStartedAt;
emitProgress("complete");
return writtenBytes;
} finally {
reader.releaseLock();
}
Expand Down
37 changes: 36 additions & 1 deletion packages/ducjs/tests/serialize-parse.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,10 @@ import { join } from "node:path";
import { gunzipSync, gzipSync } from "fflate";

import * as ducjs from "../src";
import { importDucStreamToOpfs } from "../src/opfs-import";
import {
importDucStreamToOpfs,
type DucStreamImportProgress,
} from "../src/opfs-import";
import { parseDuc } from "../src/parse";

describe("DUC streaming API", () => {
Expand Down Expand Up @@ -174,6 +177,38 @@ describe("DUC streaming API", () => {
expect(restored).toEqual(payload);
});

test("reports cumulative stream wait and processing progress", async () => {
const payload = new Uint8Array(16 + 257);
payload.set(new TextEncoder().encode("SQLite format 3\0"));
const compressed = gzipSync(payload);
const progress: DucStreamImportProgress[] = [];
const stream = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(compressed);
controller.close();
},
});

await importDucStreamToOpfs(
stream,
{ writeChunk() {} },
{
writeBufferSize: 32,
progressIntervalMs: Number.MAX_SAFE_INTEGER,
onProgress: (event) => progress.push(event),
},
);

expect(progress).toHaveLength(1);
expect(progress[0]).toMatchObject({
phase: "complete",
compressedBytes: compressed.byteLength,
uncompressedBytes: payload.byteLength,
});
expect(progress[0].readWaitMs).toBeGreaterThanOrEqual(0);
expect(progress[0].processingMs).toBeGreaterThanOrEqual(progress[0].opfsWriteMs);
});

test("copies raw SQLite into bounded OPFS importer writes", async () => {
const payload = new Uint8Array(16 + 257);
payload.set(new TextEncoder().encode("SQLite format 3\0"));
Expand Down
Loading