diff --git a/docs/resumable-downloads-plan.md b/docs/resumable-downloads-plan.md new file mode 100644 index 0000000..f52d31e --- /dev/null +++ b/docs/resumable-downloads-plan.md @@ -0,0 +1,138 @@ +# Resumable save / download — plan + +Status: **implemented 2026-10-03** (items 1–7; item 6 without the wake lock). +Verified end to end in Chromium through a connection-dropping proxy: a 19 MB +save paused after its retries with 8 MiB kept, resumed by itself after a +reload from byte 8388608, and the stored file's SHA-256 matched the server's. +As built: prepare is `GET /api/download/:id/prepare` (states ready / working / +legacy / failed); the partial's owner is a `.part.json` sidecar in OPFS +(no IndexedDB record); the queue of unfinished saves is `localStorage.ytpSaveQueue` +(`SaveQueue` in app.js); pure helpers live in `frontend/resume-core.js`. Goal: saving a long video (≥ 1 h, hundreds of MB) +to the device must survive dropped connections, app backgrounding, reloads and +the slow homelab link — continuing from where it stopped instead of starting over. + +## Why long saves get stuck today + +| Piece | Today | Problem for big files | +|---|---|---| +| `frontend/app.js` `opfsDownload()` (~273) | one `fetch('/api/download/')` for the whole file | any drop = total loss; no progress survives a reload | +| `frontend/opfs.js` `writeFromResponse()` (~178), `opfs-worker.js` (~40) | writes a `.part`, checks `Content-Length`, **deletes the `.part` on any failure** | the bytes already received are thrown away | +| `server/server.js` `/api/download/:id` (~1549) → `cachedDownloadResponse()` (~1527) | streams the cached copy with `Content-Length` only | **no `Range` / `Accept-Ranges`** → the client cannot ask for "the rest" | +| `/api/download/:id` when the server has no copy yet (~1620) | streams yt-dlp output live | not resumable at all; one stall kills it | +| iOS / Safari | backgrounding suspends `fetch` | the single long request dies whenever the phone locks | + +Already in place and reusable: `rangeFileResponse()` (~1265, used by +`/api/media/:id` and `/api/export/:id`) does correct 206/`Content-Range`; +`/api/media/:id/status` reports cache-job progress; `sha256.js` hashes +incrementally; the OPFS worker uses `createSyncAccessHandle`, which can write at +any offset; Bun `idleTimeout` is already 0. + +## Design + +**Rule: the device only ever downloads a finished, validated server copy, in +byte ranges.** Fetching from YouTube is the server's job (the media cache), +never part of a device transfer. + +``` +tap Save ──► POST /api/media/:id/prepare ──► server media-cache job (HIGH) + │ │ status/progress + │◄── poll GET /api/media/:id/status ◄────────┘ "Preparing on server 42%" + ▼ ready: { gen, size, sha256 } + for each missing chunk (8 MiB): + GET /api/media/:id?g= Range: bytes=a-b If-Range: "." + write at offset a (OPFS sync handle, in the worker) + record bytesDone in IndexedDB + all chunks ──► hash whole .part, compare sha256 ──► rename to final +``` + +### Server (Bun) — `server/server.js`, `server/media-cache.js` + +1. **Range on the download path.** `/api/download/:id` with a cached copy (and + uploads) answers through `rangeFileResponse()`, plus `Accept-Ranges: bytes` + and a strong `ETag: "."` (uploads: `""`). Honor `If-Range`: when + the ETag no longer matches (copy was re-downloaded → new gen), send a full 200 + so the client knows to restart. Keep `X-Content-SHA256`. +2. **Generation-pinned URLs.** Chunks use `/api/media/:id?g=` (already + immutable + ranged). A finished copy is never rewritten in place (new gen = + new file), so bytes for one gen never change under the client. +3. **Pin while downloading.** Eviction must not remove a copy a device is in the + middle of fetching: `media.touch(id)` on every ranged hit already refreshes + LRU; add a short "in transfer" protect window (reuse `EVICT_PROTECT_MS`). +4. **Prepare endpoint.** `POST /api/media/:id/prepare` → `ensureCached(id, {priority: HIGH})` + and return `status()` immediately (no waiting). `status()` gains + `{ progress, gen, size, sha256 }` so the client can show + "Preparing on server" with a percentage and knows the exact target. +5. **Long videos.** `MAX_SAVE_SECONDS` / `MEDIA_AUTO_MAX_SECONDS` stay as the + policy limits; when they refuse, `prepare` returns the reason so the UI can + say so instead of hanging. The yt-dlp live-stream path of `/api/download` + stays only as a legacy fallback for short videos. +6. **Uploads.** Same Range/ETag treatment (they already sit on disk; with the + USB drive move, the backup-dir fallback in `uploads.js` serves either copy — + both are byte-identical, so the ETag is the same). + +### Browser — `frontend/app.js`, `frontend/opfs-worker.js`, `frontend/opfs.js`, `frontend/device-db.js` + +1. **Transfer record (IndexedDB, `device-db.js`).** One row per saving video: + `{ id, gen, size, sha256, chunk, bytesDone, state: preparing|downloading|verifying|done|failed, updatedAt, error }`. + This is what survives reloads, crashes and app kills. +2. **Chunked worker download (`opfs-worker.js`).** Replace the single fetch with a + loop over missing ranges: `Range: bytes=start-end`, `If-Range: "."`, + per-chunk `AbortController` timeout (60 s), up to 5 retries with backoff + (2 s → 30 s). Write with `accessHandle.write(buf, { at: start })`, `flush()` + every chunk, update `bytesDone` after the flush (so the record never claims + bytes that are not on disk). On start, trust `min(record.bytesDone, .part size)`. + - 200 instead of 206 or an ETag mismatch → the server copy changed: truncate + the `.part`, reset the record to the new gen, start over (rare). + - 416 → `.part` is longer than the file: truncate to `size` and verify. +3. **Never delete progress on failure.** `.part` + record are kept on errors and + only removed on user cancel, on success, or when the record is older than 7 + days (sweep at startup, alongside OPFS quota checks). +4. **Verify at the end, not during.** Hashing across resumes is done by + re-reading the finished `.part` in the worker with `sha256.js` (no hash state + to persist), then compare with the server's `sha256`; mismatch → discard and + restart once, then mark failed. Rename `.part` → final as today. +5. **Auto-resume triggers.** App start, `online`, `visibilitychange` → visible, + and the SW `sync` event where supported: resume every record in + `downloading`/`preparing`. One transfer at a time (the homelab uplink is the + bottleneck), next in queue starts when one finishes. +6. **Browsers without worker sync access handles.** Detect up front. Where only + `createWritable` exists (some desktop Chromium contexts), use + `createWritable({ keepExistingData: true })` + `seek(start)` per chunk. If + neither exists, keep today's "not supported" error — do not fall back to an + unresumable path silently. +7. **UI.** + - Save button / Downloads list show `Preparing on server 40%`, + `Downloading 312 / 742 MB`, `Paused — will resume`, `Verifying…`. + - Pause / Resume / Cancel per item in **Downloads** (Cancel deletes `.part`). + - A paused item resumes by itself on the triggers above; a toast only on + final success or a hard failure (with the reason from the server). + - Keep the screen awake (existing wake-lock helper) while a foreground + download runs, opt-out in Settings, so iOS does not suspend it. +8. **Save-to-device export** (`exportToDevice`, ~8000) already uses the ranged + `/api/export/:id`; no change beyond using the same ETag. + +### Browser support matrix (target) + +| Browser | Write at offset | Resume across reload | Notes | +|---|---|---|---| +| Chrome / Edge / Android Chrome | worker sync handle | yes | | +| Safari / iOS 16.4+ (PWA + tab) | worker sync handle | yes | suspended when backgrounded → resumes on `visibilitychange` | +| Firefox 111+ | worker sync handle | yes | | +| Older browsers without OPFS sync handles | — | — | clear "not supported" message, Save-to-device export still works | + +## Work items (each one commit, tests first) + +1. Server: Range + strong ETag + `If-Range` on `/api/download/:id` (cached + uploads). Tests: 206 slices, full 200 on ETag mismatch, 416 past the end. +2. Server: `POST /api/media/:id/prepare`; `status()` returns `{ progress, gen, size, sha256 }`; in-transfer eviction protect. Tests in `media-cache.test.js`. +3. Browser: transfer record store in `device-db.js` + startup sweep (node tests with a fake IDB). +4. Browser: chunked, resumable worker download in `opfs-worker.js` (pure chunk planner + retry policy as testable functions; node tests with a fake fetch that drops mid-chunk, returns 200 on ETag change, and 416). +5. Browser: wire `opfsDownload()` / `preload()` to prepare → poll → chunked download; auto-resume triggers; one-at-a-time queue. +6. UI: progress states, Pause / Resume / Cancel in Downloads, wake lock while downloading. +7. End-to-end check against a local server with a 1 h+ fixture: kill the network mid-way (Playwright `context.setOffline`), reload the page, confirm it resumes from the last chunk and the final SHA-256 matches; run the same in WebKit (Windows Playwright, see CLAUDE.md) for Safari behaviour. + +## Out of scope + +- Resuming the server-side YouTube fetch itself (the media cache already restarts + jobs on boot and retries with backoff). +- Background downloads while the iOS app is fully closed (no Background Fetch on + iOS); the transfer resumes the next time the app is opened. diff --git a/frontend/app.js b/frontend/app.js index 89f2c8b..05863df 100755 --- a/frontend/app.js +++ b/frontend/app.js @@ -270,7 +270,127 @@ const SearchLibrary = (() => { // OPFS bridge wrappers — return the same shape as the Tauri cache_* commands // so the rest of app.js works without changes. -async function opfsDownload(videoId, { mux = false } = {}) { +// ---------- Save queue: resumable saves survive reloads and dropped links ---------- +// Every save is recorded here when it starts and removed when it finishes or +// fails for good. A save that only paused (partial kept on the device) stays +// listed and is resumed when the app opens, comes back online or returns to +// the foreground — and every 2 minutes while the app is visible. +const SaveQueue = (() => { + const KEY = 'ytpSaveQueue'; + const MAX_AGE = 7 * 24 * 3600 * 1000; + const prog = new Map(); + let running = false; + const load = () => { try { return JSON.parse(localStorage.getItem(KEY) || '{}') || {}; } catch { return {}; } }; + const store = (q) => { try { localStorage.setItem(KEY, JSON.stringify(q)); } catch { /* storage blocked */ } }; + + function add(v) { + if (!v || !v.id || v.custom) return; + const q = load(); + q[v.id] = { id: v.id, title: v.title || '', channel: v.channel || '', thumbnail: v.thumbnail || '', + duration: v.duration || 0, channelUrl: v.channelUrl || '', channelId: v.channelId || '', + at: (q[v.id] && q[v.id].at) || Date.now() }; + store(q); + } + function remove(id) { const q = load(); if (q[id]) { delete q[id]; store(q); } } + const pending = () => Object.values(load()); + const paused = () => pending().filter((v) => !downloading.has(v.id) && !cachedIds.has(v.id)); + const get = (id) => prog.get(id) || null; + + function label(p) { + if (!p) return '⏳ Saving for offline…'; + const mb = (n) => (n / 1048576).toFixed(n >= 1048576 * 100 ? 0 : 1); + if (p.phase === 'preparing') return `⏳ Preparing on server… ${Math.round((p.elapsed || 0) / 1000)} s`; + if (p.phase === 'verifying') return '✓ Verifying…'; + if (p.total) return `⬇ ${mb(p.received)} / ${mb(p.total)} MB`; + return '⏳ Saving for offline…'; + } + function paint(id) { + const row = document.querySelector(`.card.downloading[data-id="${CSS.escape(id)}"]`); + if (!row) return; + const p = get(id); + const ch = row.querySelector('.card-channel'); + if (ch) ch.textContent = label(p); + const bar = row.querySelector('.dl-bar'); + if (bar) { + const pct = p && p.total ? Math.max(1, Math.min(100, (p.received / p.total) * 100)) : null; + bar.classList.toggle('determinate', pct !== null); + bar.style.width = pct !== null ? pct + '%' : ''; + } + } + function progress(id, p) { if (p) prog.set(id, p); else prog.delete(id); paint(id); } + + async function resumeAll() { + if (running || !navigator.onLine) return; + running = true; + try { + for (const v of pending()) { + if (cachedIds.has(v.id)) { remove(v.id); continue; } + if (downloading.has(v.id)) continue; + if (Date.now() - v.at > MAX_AGE) { cancel(v.id); continue; } + await preload(v, { quiet: true }); // one at a time: the homelab link is the bottleneck + if (!navigator.onLine) break; + } + } finally { running = false; } + if (view.type === 'downloads') renderList(); + } + function cancel(id) { + remove(id); + prog.delete(id); + if (window.OPFS && window.OPFS.discardPartial) window.OPFS.discardPartial(id); + } + function start() { + setTimeout(resumeAll, 4000); + window.addEventListener('online', () => resumeAll()); + document.addEventListener('visibilitychange', () => { if (document.visibilityState === 'visible') resumeAll(); }); + setInterval(() => { if (document.visibilityState === 'visible' && paused().length) resumeAll(); }, 120000); + } + return { add, remove, pending, paused, progress, get, label, resumeAll, cancel, start }; +})(); + +// Resumable save (docs/resumable-downloads-plan.md): ask the server to get +// the video first (never blocks), then fetch its finished copy in byte ranges +// inside the worker. Returns null when this video must use the legacy +// single-request save (HEVC copy this device can't play, too long for the +// server cache, no worker/OPFS sync access), otherwise the save result — +// { ok:false, paused:true } means the partial is kept and will resume. +async function resumableOpfsSave(videoId, qs, hevc, onProgress) { + if (typeof window.OPFS.downloadVideo !== 'function' || typeof Worker === 'undefined') return null; + const enc = encodeURIComponent(videoId); + const prepUrl = `/api/download/${enc}/prepare${hevc ? '?hevc=1' : ''}`; + const dlUrl = `/api/download/${enc}${qs ? '?' + qs : ''}`; + const report = (info) => { if (onProgress) { try { onProgress(info); } catch { /* UI only */ } } }; + const started = Date.now(); + for (let round = 0; round < 2; round++) { + let prep = null; + for (;;) { + let j; + try { + const r = await fetch(prepUrl, { cache: 'no-store' }); + j = await r.json(); + } catch (err) { + return { ok: false, paused: true, error: 'connection lost' }; + } + if (j.state === 'ready') { prep = j; break; } + if (j.state === 'legacy') return null; + if (!j.ok || j.state === 'failed') return { ok: false, error: j.reason || j.error || 'the server could not get this video' }; + report({ phase: 'preparing', elapsed: Date.now() - started }); + if (Date.now() - started > 3 * 3600 * 1000) return { ok: false, paused: true, error: 'the server is still preparing it' }; + await new Promise((r) => setTimeout(r, 3000)); + } + const w = await window.OPFS.downloadVideo(videoId, dlUrl, { + resumable: { etag: prep.etag, size: prep.size, sha256: prep.sha256 || null, ext: prep.ext || 'mp4' }, + onProgress: (received, total, m) => report({ phase: m && m.verifying ? 'verifying' : 'downloading', received, total }), + }); + if (w.ok) return { ok: true, cached: true, sha256: w.sha256 || null, expectedSha: w.expectedSha || null, size: w.size || 0 }; + if (w.changed) continue; // the server copy was replaced — start over on the new one + if (w.paused) return { ok: false, paused: true, error: w.error || 'connection lost', received: w.received || 0 }; + if (w.fallback) return null; + return { ok: false, error: w.error || 'save failed' }; + } + return { ok: false, error: 'the server copy kept changing — try again' }; +} + +async function opfsDownload(videoId, { mux = false, onProgress = null } = {}) { if (!window.OPFS || !window.OPFS.isSupported()) { return { ok: false, error: 'OPFS not supported in this browser' }; } @@ -285,6 +405,9 @@ async function opfsDownload(videoId, { mux = false } = {}) { const qs = params.toString(); const url = `/api/download/${encodeURIComponent(videoId)}${qs ? '?' + qs : ''}`; + const resumed = await resumableOpfsSave(videoId, qs, params.get('hevc') === '1', onProgress); + if (resumed) return resumed; + // Preferred path: a dedicated Web Worker does the fetch AND the OPFS writes, // so a big save never touches the main thread (no UI jank, no audio // stutter). ANY worker failure — unsupported API or a mid-download error — @@ -1428,20 +1551,29 @@ async function preload(video, { quiet = false, mux = false } = {}) { // finishes in the `finally` below. if (current && current.meta && current.meta.id === id) updateNowPlayingActions(); if (!quiet) toast(`Saving “${video.title}” for offline…`); + SaveQueue.add(video); try { - const res = await API.cacheDownload(id, { mux }); + const res = await API.cacheDownload(id, { mux, onProgress: (p) => SaveQueue.progress(id, p) }); + if (res && res.paused) { + // The partial stays on the device; SaveQueue picks it up again. + if (!quiet) toast(`Paused “${video.title}” — it will continue when the connection is back`); + } else { + SaveQueue.remove(id); + } if (res && res.ok && res.cached) { cachedIds.add(id); cacheMutations++; recordDeviceFile(id, res); warmThumb(thumbUrlFor(id, video)); if (!quiet) toast(`Saved “${video.title}” ✓`); - } else if (!quiet) { + } else if (!quiet && !(res && res.paused)) { toast('⚠ ' + ((res && res.error) || 'Could not save video')); } } catch (e) { + SaveQueue.remove(id); if (!quiet) toast('⚠ Saving not supported in this build.'); } finally { + SaveQueue.progress(id, null); downloading.delete(id); downloadMeta.delete(id); markCardCacheState(id, cachedIds.has(id) ? 'cached' : 'none'); @@ -8311,7 +8443,8 @@ function renderDownloads() { c.innerHTML = ''; const active = [...downloadMeta.values()]; - if (!active.length) { + const paused = SaveQueue.paused(); + if (!active.length && !paused.length) { const empty = document.createElement('div'); empty.className = 'empty-state'; empty.innerHTML = ` @@ -8329,25 +8462,73 @@ function renderDownloads() { const note = document.createElement('div'); note.className = 'saved-summary'; - note.innerHTML = `${active.length}download${active.length === 1 ? '' : 's'} in progress`; + const n = active.length + paused.length; + note.innerHTML = `${n}download${n === 1 ? '' : 's'}${paused.length ? ` · ${paused.length} paused` : ' in progress'}`; c.appendChild(note); - active.forEach((v) => { + const rowFor = (v, isPaused) => { const row = document.createElement('div'); - row.className = 'card downloading'; + row.className = 'card downloading' + (isPaused ? ' dl-paused' : ''); row.dataset.id = v.id; row.innerHTML = `
- +
-
⏳ Saving for offline…
+
`; row.querySelector('.card-title').textContent = v.title || v.id; + return row; + }; + + active.forEach((v) => { + c.appendChild(rowFor(v, false)); + // paint the live label/bar now (the worker keeps updating it) + SaveQueue.progress(v.id, SaveQueue.get(v.id)); + }); + + paused.forEach((v) => { + const row = rowFor(v, true); + row.querySelector('.card-channel').textContent = '⏸ Paused — continues when the connection is back'; + const acts = document.createElement('div'); + acts.className = 'dl-actions'; + const resume = document.createElement('button'); + resume.type = 'button'; + resume.className = 'dl-act'; + resume.textContent = '▶ Resume'; + resume.onclick = (e) => { e.stopPropagation(); preload(v, { quiet: true }); renderList(); }; + const cancel = document.createElement('button'); + cancel.type = 'button'; + cancel.className = 'dl-act danger'; + cancel.textContent = 'Cancel'; + cancel.onclick = (e) => { + e.stopPropagation(); + showModal(`Cancel saving “${v.title || v.id}”?`, document.createTextNode('The part already downloaded will be deleted.'), [ + { label: 'Keep', onClick: closeModal }, + { label: 'Cancel save', danger: true, onClick: () => { SaveQueue.cancel(v.id); closeModal(); renderList(); } }, + ]); + }; + acts.append(resume, cancel); + row.querySelector('.card-info').appendChild(acts); c.appendChild(row); }); + + // Fill in how far each paused save got (read from the partial on disk). + if (paused.length && window.OPFS && window.OPFS.listPartials) { + window.OPFS.listPartials().then((parts) => { + for (const p of parts) { + const row = c.querySelector(`.card.dl-paused[data-id="${CSS.escape(p.id)}"]`); + if (!row || !p.size) continue; + const bar = row.querySelector('.dl-bar'); + bar.classList.add('determinate'); + bar.style.width = Math.max(1, Math.min(100, (p.received / p.size) * 100)) + '%'; + row.querySelector('.card-channel').textContent = + `⏸ Paused at ${(p.received / 1048576).toFixed(0)} / ${(p.size / 1048576).toFixed(0)} MB — continues when the connection is back`; + } + }).catch(() => {}); + } } // ============================================================================ @@ -11125,6 +11306,8 @@ async function boot() { } data.playlists.forEach(preloadPlaylist); preloadPinnedPlaylists(); + // Continue saves that were interrupted (closed app, dropped connection). + if (WEB) SaveQueue.start(); // Not awaited — artwork backfill must never delay first paint. warmOfflineThumbs(); diff --git a/frontend/opfs-worker.js b/frontend/opfs-worker.js index 36b716b..d854d21 100644 --- a/frontend/opfs-worker.js +++ b/frontend/opfs-worker.js @@ -9,7 +9,12 @@ * 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 } + * 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 @@ -18,6 +23,10 @@ * 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 * ========================================================================== */ @@ -26,6 +35,7 @@ // 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. try { importScripts('/sha256.js'); } catch { /* hashing unavailable — save still works */ } +importScripts('/resume-core.js'); async function getVideosDir() { const root = await navigator.storage.getDirectory(); @@ -37,8 +47,146 @@ function extFromContentType(ct) { 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 }); +} + self.onmessage = async (e) => { - const { videoId, url } = e.data || {}; + const { videoId, url, resumable } = e.data || {}; if ( typeof navigator === 'undefined' || @@ -51,6 +199,12 @@ self.onmessage = async (e) => { 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 { @@ -100,34 +254,7 @@ self.onmessage = async (e) => { throw new Error('integrity check failed (content hash mismatch)'); } - // Finalize: .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. - try { await dir.removeEntry(filename); } catch { /* no previous copy */ } - let renamed = false; - if (typeof partHandle.move === 'function') { - try { await partHandle.move(filename); renamed = true; } catch { /* copy below */ } - } - if (!renamed) { - 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); - } + await finalizePart(dir, partHandle, partName, filename); self.postMessage({ type: 'done', ext, sha256, expectedSha, size: offset }); } catch (err) { diff --git a/frontend/opfs.js b/frontend/opfs.js index 7580c8e..be24596 100644 --- a/frontend/opfs.js +++ b/frontend/opfs.js @@ -145,7 +145,9 @@ // workers. Resolves { ok:true } on success, { ok:false, error } on a real // failure, or { ok:false, fallback:true } when the worker path is // unavailable and the caller should use writeFromResponse instead. - downloadVideo(videoId, url) { + // opts.resumable = { etag, size, sha256, ext } switches the worker to the + // ranged, resumable save; opts.onProgress(received, total, info) reports. + downloadVideo(videoId, url, opts = {}) { return new Promise((resolve) => { let worker; try { @@ -163,10 +165,14 @@ if (m.type === 'done') finish({ ok: true, sha256: m.sha256 || null, expectedSha: m.expectedSha || null, size: m.size || 0 }); else if (m.type === 'unsupported') finish({ ok: false, fallback: true }); else if (m.type === 'error') finish({ ok: false, error: m.error }); - // 'progress' messages are informational; ignored here + else if (m.type === 'paused') finish({ ok: false, paused: true, received: m.received || 0, error: m.error }); + else if (m.type === 'changed') finish({ ok: false, changed: true }); + else if (m.type === 'progress' && opts.onProgress) { + try { opts.onProgress(m.received || 0, m.total || 0, m); } catch { /* UI only */ } + } }; worker.onerror = () => finish({ ok: false, fallback: true }); - worker.postMessage({ videoId, url }); + worker.postMessage({ videoId, url, resumable: opts.resumable || null }); }); }, @@ -249,8 +255,8 @@ try { for await (const [name, handle] of dir.entries()) { if (handle.kind !== 'file') continue; - // Skip .part temporary files - if (name.endsWith('.part')) continue; + // Skip .part temporary files and resumable-save sidecars + if (name.endsWith('.part') || name.endsWith('.part.json')) continue; const dot = name.lastIndexOf('.'); const id = dot > -1 ? name.slice(0, dot) : name; const file = await handle.getFile(); @@ -260,6 +266,37 @@ return items; }, + // Paused resumable saves: [{ id, received, size }] from ..part + // and its .part.json sidecar (written by opfs-worker.js). + async listPartials() { + const out = []; + try { + const dir = await getRoot(); + for await (const [name, handle] of dir.entries()) { + if (handle.kind !== 'file' || !name.endsWith('.part.json')) continue; + const partName = name.slice(0, -5); + let meta = null, received = 0; + try { meta = JSON.parse(await (await handle.getFile()).text()); } catch { /* unreadable */ } + try { received = (await (await dir.getFileHandle(partName)).getFile()).size; } catch { /* no data yet */ } + const id = partName.replace(/\.[^.]+\.part$/, ''); + out.push({ id, received, size: (meta && meta.size) || 0 }); + } + } catch { /* OPFS not available */ } + return out; + }, + + // Drop a paused save's partial data (Cancel in Downloads). + async discardPartial(videoId) { + try { + const dir = await getRoot(); + const names = []; + for await (const [name] of dir.entries()) { + if ((name.endsWith('.part') || name.endsWith('.part.json')) && name.startsWith(videoId + '.')) names.push(name); + } + for (const name of names) await dir.removeEntry(name).catch(() => {}); + } catch { /* OPFS not available */ } + }, + // Total bytes stored async totalSize() { const items = await this.listVideos(); diff --git a/frontend/resume-core.js b/frontend/resume-core.js new file mode 100644 index 0000000..3408efb --- /dev/null +++ b/frontend/resume-core.js @@ -0,0 +1,44 @@ +/* ============================================================================ + * resume-core — pure helpers for resumable saves (docs/resumable-downloads-plan.md) + * + * Framework-free: a plain