Files
J621/frontend/src/features/optimize/stream.ts
T
JakeBreath e7657a3dd0 Fix the OPFS output lifecycle (locked stream on close)
StreamTarget closes the underlying writer when the output is finalized —
that is also what commits an OPFS file — so closing it ourselves threw
"Can not close locked stream". finalize() now simply reads the committed
file back, and every failure path (invalid attempt, encode error,
validation failure) aborts the write and deletes the temporary entry.
2026-09-17 21:01:15 -05:00

149 lines
4.1 KiB
TypeScript

import { StreamTarget, type StreamTargetChunk } from "mediabunny";
import {
createStoredFile,
deleteStoredFile,
opfsSupported,
purgeStaleStoredFiles,
} from "./opfs";
interface CollectedChunk {
blob: Blob;
position: number;
sequence: number;
}
export interface StreamOutput {
target: StreamTarget;
finalize(type: string): Promise<{ blob: Blob; storageKey?: string }>;
abort(): Promise<void>;
}
/**
* Stream the muxed output somewhere that is not one giant ArrayBuffer — a 4K
* file would otherwise blow the tab's memory.
*
* Preferred: a real Origin Private File System file with random access, which
* is what the muxer expects (it rewrites earlier byte ranges for headers).
* StreamTarget closes the underlying writer when the output is finalized,
* which also commits the OPFS file. Fallback: collect chunks and assemble
* them, with the newest write winning per byte range.
*/
export async function createStreamOutput(
extension: string,
): Promise<StreamOutput> {
if (opfsSupported()) {
const stored = await createStoredFile(extension);
if (stored) {
void purgeStaleStoredFiles();
const writable = await stored.handle.createWritable();
let finalized = false;
return {
target: new StreamTarget(writable),
finalize: async (type: string) => {
finalized = true;
// The writer is already closed and the file committed; just read it.
const file = await stored.handle.getFile();
return {
blob: new File([file], stored.name, { type }),
storageKey: stored.name,
};
},
abort: async () => {
if (!finalized) {
try {
await writable.abort();
} catch {
// Locked or already closed — the entry removal still applies.
}
}
await deleteStoredFile(stored.name);
},
};
}
}
return createBufferOutput();
}
function createBufferOutput(): StreamOutput {
const chunks: CollectedChunk[] = [];
let sequence = 0;
const writable = new WritableStream<StreamTargetChunk>({
write(chunk) {
chunks.push({
blob: new Blob([chunk.data.slice()]),
position: chunk.position,
sequence: sequence++,
});
},
});
return {
target: new StreamTarget(writable, {
chunked: true,
chunkSize: 8 * 1024 * 1024,
}),
finalize: async (type: string) => ({
blob: assembleBlob(chunks, type),
}),
abort: async () => {
chunks.length = 0;
},
};
}
export function assembleBlob(chunks: CollectedChunk[], type: string): Blob {
if (chunks.length === 0) {
throw new Error("Encoding produced no output.");
}
// Apply writes in arrival order — the newest write wins per byte range —
// and only then lay the segments out by position.
const ordered = [...chunks].sort((a, b) => a.sequence - b.sequence);
interface Segment {
start: number;
end: number;
sequence: number;
blob: Blob;
}
let segments: Segment[] = [];
for (const chunk of ordered) {
const start = chunk.position;
const end = start + chunk.blob.size;
const next: Segment[] = [];
for (const segment of segments) {
if (segment.end <= start || segment.start >= end) {
next.push(segment);
continue;
}
if (segment.start < start) {
next.push({
...segment,
end: start,
blob: segment.blob.slice(0, start - segment.start),
});
}
if (segment.end > end) {
next.push({
...segment,
start: end,
blob: segment.blob.slice(end - segment.start),
});
}
}
next.push({ start, end, sequence: chunk.sequence, blob: chunk.blob });
segments = next;
}
segments.sort((a, b) => a.start - b.start);
const parts: Blob[] = [];
let cursor = 0;
for (const segment of segments) {
if (segment.start !== cursor) {
throw new Error("The muxer produced an incomplete file.");
}
parts.push(segment.blob);
cursor = segment.end;
}
return new Blob(parts, { type });
}