/* ============================================================================ * opfs-worker.js — off-main-thread video download → OPFS * * One dedicated Worker per download (spawned by OPFS.downloadVideo in * opfs.js, terminated when finished). The worker does the whole job itself — * fetch from /api/download plus streaming writes via createSyncAccessHandle — * so a multi-hundred-MB save never allocates buffers or runs stream pumps on * the main thread. createSyncAccessHandle is worker-only but has wider * support than createWritable (Safari 15.2+ vs 18.2+), which also removes * the whole-file ArrayBuffer fallback the main-thread path needs on WebKit. * * In message: { videoId, url } → one streamed request (legacy) * { videoId, url, resumable: { etag, size, sha256, ext } } * → ranged, resumable save: the .part and a small * .part.json sidecar (which server copy it belongs to) * survive errors, reloads and app kills, and the next try * continues from the bytes already on disk. * Out messages: * { type: 'unsupported' } → caller falls back to main thread * { type: 'progress', received } → bytes written so far * { type: 'done', ext, sha256, expectedSha, size } * → file stored as .; * sha256 = hash of the stored bytes * (P2P content id), expectedSha = the * server's X-Content-SHA256 or null * { type: 'paused', received, error } → resumable save stopped after its * retries; .part kept for next time * { type: 'changed' } → the server copy changed; .part * dropped, caller should re-prepare * { type: 'error', error } → failed; .part cleaned up * ========================================================================== */ 'use strict'; // Incremental SHA-256 (frontend/sha256.js): the file is hashed while it is // written, so the device knows its content id without re-reading the file. async function getVideosDir() { const root = await navigator.storage.getDirectory(); return root.getDirectoryHandle('videos', { create: true }); } function extFromContentType(ct) { ct = ct || 'video/mp4'; return ct.includes('webm') ? 'webm' : ct.includes('ogg') ? 'ogg' : 'mp4'; } // .part → permanent name. Prefer the native rename, but treat ANY move() // failure as "unavailable" and fall back to a chunked copy — WebKit's move() // has a different signature/behavior than Chrome's and throws TypeError // ("Not enough arguments") rather than being absent. async function finalizePart(dir, partHandle, partName, filename) { try { await dir.removeEntry(filename); } catch { /* no previous copy */ } if (typeof partHandle.move === 'function') { try { await partHandle.move(filename); return; } catch { /* copy below */ } } const finalHandle = await dir.getFileHandle(filename, { create: true }); const out = await finalHandle.createSyncAccessHandle(); try { const file = await partHandle.getFile(); const CHUNK = 8 * 1024 * 1024; let pos = 0; while (pos < file.size) { const buf = await file.slice(pos, pos + CHUNK).arrayBuffer(); out.write(new Uint8Array(buf), { at: pos }); pos += buf.byteLength; } out.truncate(file.size); out.flush(); } finally { out.close(); } await dir.removeEntry(partName); } const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); async function readSidecar(dir, name) { try { return JSON.parse(await (await (await dir.getFileHandle(name)).getFile()).text()); } catch { return null; } } async function writeSidecar(dir, name, obj) { const h = await dir.getFileHandle(name, { create: true }); const a = await h.createSyncAccessHandle(); try { const bytes = new TextEncoder().encode(JSON.stringify(obj)); a.truncate(0); a.write(bytes, { at: 0 }); a.flush(); } finally { a.close(); } } const MAX_TRIES = 6; // per chunk, with backoff 2 s … 30 s const IDLE_MS = 30000; // abort a request that delivers nothing for 30 s async function resumableSave(videoId, url, target) { const R = self.ResumeCore; const dir = await getVideosDir(); const filename = videoId + '.' + (target.ext || 'mp4'); const partName = filename + '.part'; const sideName = partName + '.json'; const partHandle = await dir.getFileHandle(partName, { create: true }); const sidecar = await readSidecar(dir, sideName); let access = await partHandle.createSyncAccessHandle(); let pos = 0; try { pos = R.resumeFrom(sidecar, access.getSize(), target); access.truncate(pos); access.flush(); await writeSidecar(dir, sideName, { etag: target.etag, size: target.size, sha256: target.sha256 || null }); self.postMessage({ type: 'progress', received: pos, total: target.size, resumed: pos > 0 }); let tries = 0; let lastPost = 0; for (let r = R.nextRange(pos, target.size); r; r = R.nextRange(pos, target.size)) { const ctl = new AbortController(); let idle = setTimeout(() => ctl.abort(), IDLE_MS); const bump = () => { clearTimeout(idle); idle = setTimeout(() => ctl.abort(), IDLE_MS); }; try { const res = await fetch(url, { headers: { Range: `bytes=${r.start}-${r.end}`, 'If-Range': target.etag }, signal: ctl.signal, cache: 'no-store' }); if (res.status === 200 || res.status === 416) { // The server no longer has the copy this partial belongs to. try { await res.body?.cancel(); } catch { /* fine */ } access.truncate(0); access.close(); access = null; try { await dir.removeEntry(partName); } catch { /* gone */ } try { await dir.removeEntry(sideName); } catch { /* gone */ } self.postMessage({ type: 'changed' }); return; } if (res.status !== 206) throw new Error('HTTP ' + res.status); const cr = R.parseContentRange(res.headers.get('content-range')); if (!cr || cr.start !== pos) throw new Error('unexpected range from server'); const reader = res.body.getReader(); for (;;) { const { done, value } = await reader.read(); if (done) break; bump(); access.write(value, { at: pos }); pos += value.byteLength; const now = Date.now(); if (now - lastPost > 250) { lastPost = now; self.postMessage({ type: 'progress', received: pos, total: target.size }); } } access.flush(); if (pos < r.end + 1) throw new Error('connection closed early'); tries = 0; } catch (err) { access.flush(); // keep every byte that did arrive if (++tries >= MAX_TRIES) { self.postMessage({ type: 'paused', received: pos, total: target.size, error: err && err.message ? err.message : String(err) }); return; } await sleep(R.backoffMs(tries - 1)); } finally { clearTimeout(idle); } } access.truncate(target.size); access.flush(); } finally { if (access) access.close(); } // Verify the whole file once at the end (no hash state to carry across // resumes): read it back in slices and compare with the server's hash. self.postMessage({ type: 'progress', received: target.size, total: target.size, verifying: true }); let sha256 = null; if (self.Sha256) { const h = self.Sha256.create(); const file = await partHandle.getFile(); const CHUNK = 8 * 1024 * 1024; for (let p = 0; p < file.size; p += CHUNK) h.update(new Uint8Array(await file.slice(p, p + CHUNK).arrayBuffer())); sha256 = h.hex(); } const expectedSha = /^[0-9a-f]{64}$/.test(String(target.sha256 || '')) ? target.sha256 : null; if (sha256 && expectedSha && sha256 !== expectedSha) { try { await dir.removeEntry(partName); } catch { /* gone */ } try { await dir.removeEntry(sideName); } catch { /* gone */ } throw new Error('integrity check failed (content hash mismatch)'); } await finalizePart(dir, partHandle, partName, filename); try { await dir.removeEntry(sideName); } catch { /* gone */ } self.postMessage({ type: 'done', ext: target.ext || 'mp4', sha256, expectedSha, size: target.size }); } // Write a whole file into an OPFS sub-directory (sync access handle: works on // iOS Safari, where main-thread createWritable is missing). Used for the // iPhone EQ/levelling renders (EqRender in app.js). async function writeFile(dirName, name, buffer) { const root = await navigator.storage.getDirectory(); const dir = await root.getDirectoryHandle(dirName, { create: true }); const h = await dir.getFileHandle(name, { create: true }); const a = await h.createSyncAccessHandle(); try { const bytes = new Uint8Array(buffer); a.truncate(0); a.write(bytes, { at: 0 }); a.truncate(bytes.byteLength); a.flush(); } finally { a.close(); } } self.onmessage = async (e) => { self.__ASSET_URLS__=e.data?.assetUrls || {}; try { if(!self.Sha256)importScripts(self.__ASSET_URLS__.sha256 || '/sha256.js'); } catch { /* hashing unavailable — legacy save still works */ } try { if(!self.ResumeCore)importScripts(self.__ASSET_URLS__.resume || '/resume-core.js'); } catch(error) { self.postMessage({type:'error',error:error.message}); return; } if (e.data && e.data.op === 'write') { try { await writeFile(e.data.dir, e.data.name, e.data.buffer); self.postMessage({ type: 'written' }); } catch (err) { self.postMessage({ type: 'error', error: err && err.message ? err.message : String(err) }); } return; } const { videoId, url, resumable } = e.data || {}; if ( typeof navigator === 'undefined' || !navigator.storage || typeof navigator.storage.getDirectory !== 'function' || typeof FileSystemFileHandle === 'undefined' || typeof FileSystemFileHandle.prototype.createSyncAccessHandle !== 'function' ) { self.postMessage({ type: 'unsupported' }); return; } if (resumable) { try { await resumableSave(videoId, url, resumable); } catch (err) { self.postMessage({ type: 'error', error: err && err.message ? err.message : String(err) }); } return; } let dir = null; let partName = null; try { const res = await fetch(url); if (!res.ok) { let msg = 'HTTP ' + res.status; try { msg = (await res.json()).error || msg; } catch { /* non-JSON */ } throw new Error(msg); } const ext = extFromContentType(res.headers.get('content-type')); const filename = videoId + '.' + ext; partName = filename + '.part'; dir = await getVideosDir(); const partHandle = await dir.getFileHandle(partName, { create: true }); const access = await partHandle.createSyncAccessHandle(); const hasher = self.Sha256 ? self.Sha256.create() : null; let offset = 0; try { const reader = res.body.getReader(); for (;;) { const { done, value } = await reader.read(); if (done) break; access.write(value, { at: offset }); if (hasher) hasher.update(value); offset += value.byteLength; self.postMessage({ type: 'progress', received: offset }); } access.truncate(offset); access.flush(); } finally { access.close(); } // A connection that ends early can close the body cleanly; never commit // a truncated file, it would be badged "cached" yet fail to play. const expected = Number(res.headers.get('content-length')); if (expected > 0 && offset !== expected) { throw new Error(`download cut short (${offset} of ${expected} bytes)`); } // The server names the hash of what it sent (media cache copies). A // mismatch means the bytes were damaged on the way — never keep them. const sha256 = hasher ? hasher.hex() : null; const sent = (res.headers.get('x-content-sha256') || '').trim().toLowerCase(); const expectedSha = /^[0-9a-f]{64}$/.test(sent) ? sent : null; if (sha256 && expectedSha && sha256 !== expectedSha) { throw new Error('integrity check failed (content hash mismatch)'); } await finalizePart(dir, partHandle, partName, filename); self.postMessage({ type: 'done', ext, sha256, expectedSha, size: offset }); } catch (err) { // Never leave a corrupt partial behind try { if (dir && partName) await dir.removeEntry(partName); } catch { /* gone */ } self.postMessage({ type: 'error', error: err && err.message ? err.message : String(err) }); } };