Files
ytplayer/frontend/p2p-recv-worker.js

91 lines
3.7 KiB
JavaScript

/* ============================================================================
* p2p-recv-worker.js — writes a file arriving from another device into OPFS
* while hashing it; commits it ONLY when the SHA-256 equals the expected
* content id (docs/p2p-architecture.md flow 6).
*
* In: { op:'open', videoId } start videos/<videoId>.p2p.part
* { op:'chunk', buf } ArrayBuffer (transferred)
* { op:'finish', cid, size } verify + rename to <videoId>.mp4
* { op:'abort' } drop the partial file
* Out: { op:'opened' } | { op:'progress', received } |
* { op:'done', ok:true, sha256, size } | { op:'done', ok:false, error }
* The ".part" suffix keeps OPFS.listVideos() from ever listing a partial file.
* ========================================================================== */
'use strict';
importScripts('/sha256.js');
let dir = null;
let handle = null;
let access = null;
let partName = '';
let finalName = '';
let hasher = null;
let offset = 0;
let lastProgress = 0;
async function cleanup() {
try { if (access) access.close(); } catch { /* closed */ }
access = null;
try { if (dir && partName) await dir.removeEntry(partName); } catch { /* gone */ }
}
// Messages are handled strictly one after another: an async handler would
// otherwise let 'chunk' or 'abort' run while 'open' is still awaiting OPFS.
let chain = Promise.resolve();
self.onmessage = (e) => { chain = chain.then(() => onOp(e.data || {})); };
async function onOp(m) {
try {
if (m.op === 'open') {
const root = await navigator.storage.getDirectory();
dir = await root.getDirectoryHandle('videos', { create: true });
partName = `${m.videoId}.p2p.part`;
finalName = `${m.videoId}.mp4`;
handle = await dir.getFileHandle(partName, { create: true });
access = await handle.createSyncAccessHandle();
access.truncate(0);
hasher = self.Sha256.create();
offset = 0;
self.postMessage({ op: 'opened' });
} else if (m.op === 'chunk') {
const u8 = new Uint8Array(m.buf);
access.write(u8, { at: offset });
hasher.update(u8);
offset += u8.byteLength;
if (offset - lastProgress > 1024 * 1024) { lastProgress = offset; self.postMessage({ op: 'progress', received: offset }); }
} else if (m.op === 'finish') {
access.truncate(offset);
access.flush();
access.close();
access = null;
const sha256 = hasher.hex();
if (offset !== m.size) throw new Error(`size mismatch (${offset} of ${m.size})`);
if (sha256 !== m.cid) throw new Error('content hash mismatch');
try { await dir.removeEntry(finalName); } catch { /* no previous copy */ }
let renamed = false;
if (typeof handle.move === 'function') { try { await handle.move(finalName); renamed = true; } catch { /* copy below */ } }
if (!renamed) {
const out = await (await dir.getFileHandle(finalName, { create: true })).createSyncAccessHandle();
try {
const file = await handle.getFile();
for (let pos = 0; pos < file.size; pos += 8 * 1024 * 1024) {
const buf = new Uint8Array(await file.slice(pos, pos + 8 * 1024 * 1024).arrayBuffer());
out.write(buf, { at: pos });
}
out.truncate(file.size);
out.flush();
} finally { out.close(); }
await dir.removeEntry(partName);
}
partName = '';
self.postMessage({ op: 'done', ok: true, sha256, size: offset });
} else if (m.op === 'abort') {
await cleanup();
self.postMessage({ op: 'done', ok: false, error: 'aborted' });
}
} catch (err) {
await cleanup();
self.postMessage({ op: 'done', ok: false, error: err && err.message ? err.message : String(err) });
}
}