Resumable saves: ranged downloads that survive dropped connections, reloads and app kills
Server
- /api/download/:id answers Range with a strong ETag ("<id>.<gen>") and honours
If-Range; a stale partial gets the whole current file (streamed, so Bun does
not re-apply the Range itself). Uploads get the same treatment.
- GET /api/download/:id/prepare never blocks: ready {gen,size,sha256,etag,ext},
working (server still fetching), legacy (HEVC / too long / cache offline),
failed. Recently refused prepares are remembered for 10 minutes.
Browser
- The OPFS worker saves in 8 MiB ranges, writes at the byte offset, flushes
each chunk, retries each chunk 6 times with backoff (30 s idle timeout) and
keeps the .part plus a .part.json sidecar naming the server copy it belongs
to. A changed copy restarts cleanly; the finished file is hashed once and
checked against the server's SHA-256.
- SaveQueue remembers unfinished saves and resumes them on start, online,
return to the foreground and every 2 minutes while visible; one at a time.
- Downloads shows live MB progress, 'Preparing on server', 'Verifying', and
paused saves with Resume and Cancel (confirmed).
- listVideos ignores the sidecars; new listPartials/discardPartial helpers.
Verified in Chromium through a connection-dropping proxy: paused at 8 MiB,
auto-resumed after a reload from byte 8388608, final SHA-256 matched.
This commit is contained in:
138
docs/resumable-downloads-plan.md
Normal file
138
docs/resumable-downloads-plan.md
Normal file
@@ -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 `<file>.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/<id>')` 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=<gen> Range: bytes=a-b If-Range: "<id>.<gen>"
|
||||
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: "<id>.<gen>"` (uploads: `"<id>"`). 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=<gen>` (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: "<id>.<gen>"`,
|
||||
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.
|
||||
201
frontend/app.js
201
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 = `<span class="saved-total">${active.length}</span><span class="saved-sub">download${active.length === 1 ? '' : 's'} in progress</span>`;
|
||||
const n = active.length + paused.length;
|
||||
note.innerHTML = `<span class="saved-total">${n}</span><span class="saved-sub">download${n === 1 ? '' : 's'}${paused.length ? ` · ${paused.length} paused` : ' in progress'}</span>`;
|
||||
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 = `
|
||||
<div class="thumb">
|
||||
<img loading="lazy" decoding="async" src="${escapeHtml(v.thumbnail || '')}" alt="" />
|
||||
<img loading="lazy" decoding="async" src="${escapeHtml(v.thumbnail || thumbUrlFor(v.id, v) || '')}" alt="" />
|
||||
<div class="dl-progress"><div class="dl-bar"></div></div>
|
||||
</div>
|
||||
<div class="card-info">
|
||||
<div class="card-title"></div>
|
||||
<div class="card-channel">⏳ Saving for offline…</div>
|
||||
<div class="card-channel"></div>
|
||||
</div>`;
|
||||
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();
|
||||
|
||||
|
||||
@@ -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
|
||||
* <file>.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) {
|
||||
|
||||
@@ -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 <id>.<ext>.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();
|
||||
|
||||
44
frontend/resume-core.js
Normal file
44
frontend/resume-core.js
Normal file
@@ -0,0 +1,44 @@
|
||||
/* ============================================================================
|
||||
* resume-core — pure helpers for resumable saves (docs/resumable-downloads-plan.md)
|
||||
*
|
||||
* Framework-free: a plain <script>/importScripts global (`ResumeCore`) in the
|
||||
* browser and worker, `require`-able by `node --test`.
|
||||
* ========================================================================== */
|
||||
(function (root) {
|
||||
'use strict';
|
||||
|
||||
const CHUNK = 8 * 1024 * 1024;
|
||||
|
||||
/** Next byte range to fetch, or null when the file is complete. */
|
||||
function nextRange(pos, size, chunk = CHUNK) {
|
||||
if (!(size > 0) || pos >= size) return null;
|
||||
return { start: pos, end: Math.min(pos + chunk, size) - 1 };
|
||||
}
|
||||
|
||||
/** "bytes 0-99/1234" → { start, end, total }; null when malformed. */
|
||||
function parseContentRange(h) {
|
||||
const m = /^bytes (\d+)-(\d+)\/(\d+|\*)$/.exec(String(h || '').trim());
|
||||
if (!m) return null;
|
||||
return { start: Number(m[1]), end: Number(m[2]), total: m[3] === '*' ? null : Number(m[3]) };
|
||||
}
|
||||
|
||||
/** Wait before retry `attempt` (0-based): 2 s, 4 s, 8 s, 16 s, 30 s cap. */
|
||||
function backoffMs(attempt) {
|
||||
return Math.min(30000, 2000 * Math.pow(2, Math.max(0, attempt)));
|
||||
}
|
||||
|
||||
/**
|
||||
* Where a stored partial may resume from. The partial is only reused when
|
||||
* it was written for the very same server copy (ETag) and is not longer
|
||||
* than that copy; otherwise start over.
|
||||
*/
|
||||
function resumeFrom(sidecar, partSize, target) {
|
||||
if (!sidecar || !target || sidecar.etag !== target.etag || sidecar.size !== target.size) return 0;
|
||||
if (!(partSize > 0)) return 0;
|
||||
return Math.min(partSize, target.size);
|
||||
}
|
||||
|
||||
const ResumeCore = { CHUNK, nextRange, parseContentRange, backoffMs, resumeFrom };
|
||||
if (typeof module !== 'undefined' && module.exports) module.exports = ResumeCore;
|
||||
else root.ResumeCore = ResumeCore;
|
||||
})(typeof self !== 'undefined' ? self : this);
|
||||
31
frontend/resume-core.test.js
Normal file
31
frontend/resume-core.test.js
Normal file
@@ -0,0 +1,31 @@
|
||||
'use strict';
|
||||
const { test } = require('node:test');
|
||||
const assert = require('node:assert');
|
||||
const R = require('./resume-core');
|
||||
|
||||
test('nextRange walks the file in chunks and stops at the end', () => {
|
||||
assert.deepStrictEqual(R.nextRange(0, 20, 8), { start: 0, end: 7 });
|
||||
assert.deepStrictEqual(R.nextRange(16, 20, 8), { start: 16, end: 19 });
|
||||
assert.strictEqual(R.nextRange(20, 20, 8), null);
|
||||
assert.strictEqual(R.nextRange(0, 0, 8), null);
|
||||
});
|
||||
|
||||
test('parseContentRange reads start, end and total', () => {
|
||||
assert.deepStrictEqual(R.parseContentRange('bytes 8-15/20'), { start: 8, end: 15, total: 20 });
|
||||
assert.deepStrictEqual(R.parseContentRange('bytes 0-1/*'), { start: 0, end: 1, total: null });
|
||||
assert.strictEqual(R.parseContentRange('bytes */20'), null);
|
||||
assert.strictEqual(R.parseContentRange(null), null);
|
||||
});
|
||||
|
||||
test('backoff grows and caps at 30 s', () => {
|
||||
assert.deepStrictEqual([0, 1, 2, 3, 4, 9].map(R.backoffMs), [2000, 4000, 8000, 16000, 30000, 30000]);
|
||||
});
|
||||
|
||||
test('a partial is reused only for the same server copy', () => {
|
||||
const t = { etag: '"a.1"', size: 100 };
|
||||
assert.strictEqual(R.resumeFrom({ etag: '"a.1"', size: 100 }, 40, t), 40);
|
||||
assert.strictEqual(R.resumeFrom({ etag: '"a.0"', size: 100 }, 40, t), 0, 'older generation');
|
||||
assert.strictEqual(R.resumeFrom({ etag: '"a.1"', size: 90 }, 40, t), 0, 'size changed');
|
||||
assert.strictEqual(R.resumeFrom(null, 40, t), 0, 'no sidecar');
|
||||
assert.strictEqual(R.resumeFrom({ etag: '"a.1"', size: 100 }, 150, t), 100, 'never past the end');
|
||||
});
|
||||
@@ -4268,3 +4268,15 @@ body.landscape-fs .player-stage { touch-action: none; } /* fullscreen (real or t
|
||||
|
||||
/* P2P: "Get it from a device" next to Retry (plan 017) */
|
||||
.peer-btn { margin-left: 8px; }
|
||||
|
||||
/* Resumable saves: a real progress bar once the size is known, paused rows */
|
||||
.dl-bar.determinate { animation: none; transform: none; border-radius: 0 2px 2px 0; transition: width 0.4s ease; }
|
||||
.card.dl-paused .dl-bar { background: var(--text-dim); }
|
||||
.dl-actions { display: flex; gap: 8px; margin-top: 8px; }
|
||||
.dl-act {
|
||||
height: 32px; padding: 0 12px; border-radius: 10px; cursor: pointer;
|
||||
border: 1px solid var(--line); background: var(--bg-2); color: var(--text);
|
||||
font: 600 12px var(--ui);
|
||||
}
|
||||
.dl-act:hover { border-color: var(--accent); }
|
||||
.dl-act.danger:hover { border-color: #ff5a4f; color: #ff5a4f; }
|
||||
|
||||
@@ -70,6 +70,7 @@ const SHELL = [
|
||||
'/lyrics-core.js',
|
||||
'/stats-core.js',
|
||||
'/sha256.js',
|
||||
'/resume-core.js',
|
||||
'/sha256-wasm.js',
|
||||
'/device-db.js',
|
||||
'/hash-worker.js',
|
||||
|
||||
@@ -648,6 +648,8 @@ export function createMediaCache({
|
||||
const row = await db.getMedia(id);
|
||||
return {
|
||||
status: job ? job.status : row ? row.status : 'none',
|
||||
gen: row ? row.gen : null,
|
||||
sha256: row && row.sha256 ? row.sha256 : null,
|
||||
size: row ? row.size : 0,
|
||||
height: row ? row.height : 0,
|
||||
vcodec: row ? row.vcodec : null,
|
||||
|
||||
@@ -42,7 +42,7 @@ import { createHash } from 'node:crypto';
|
||||
import { brotliCompressSync, constants as zlibConstants } from 'node:zlib';
|
||||
import { initDb, upsertUser, recordVideoAccess, getUserData, createProfile, getProfile, saveProfile, createSharedPlaylist, getSharedPlaylist, queueInboxPlaylist, listInbox, deleteInboxItem, countInbox,
|
||||
getMedia, upsertMedia, deleteMedia, listMedia, listMediaLru, touchMedia, mediaStats } from './db.js';
|
||||
import { createMediaCache, HIGH, LOW, validateMedia } from './media-cache.js';
|
||||
import { createMediaCache, HIGH, LOW, validateMedia, MediaSkip } from './media-cache.js';
|
||||
import * as notesDb from './db.js';
|
||||
import { registerNoteRoutes, parseLrc, sanitizeLyrics } from './notes.js';
|
||||
import { createRemoteHub } from './remote.js';
|
||||
@@ -1262,16 +1262,27 @@ function cachedStreamsPayload(videoId, row) {
|
||||
}
|
||||
|
||||
// Serve a file with byte-Range support (the <video> element seeks with it).
|
||||
function rangeFileResponse(c, path, contentType, cacheControl) {
|
||||
// `etag` (strong, quoted) lets a resuming client send If-Range: when the file
|
||||
// changed under it, the Range is ignored and the whole new file comes back
|
||||
// with 200, so stale and fresh bytes are never stitched together.
|
||||
function rangeFileResponse(c, path, contentType, cacheControl, { etag = null, headers: extra = {} } = {}) {
|
||||
const file = Bun.file(path);
|
||||
const total = file.size;
|
||||
const headers = {
|
||||
'Content-Type': contentType,
|
||||
'Accept-Ranges': 'bytes',
|
||||
'Cache-Control': cacheControl,
|
||||
...(etag ? { ETag: etag } : {}),
|
||||
...extra,
|
||||
};
|
||||
const range = c.req.header('range');
|
||||
if (!range) return new Response(file, { status: 200, headers });
|
||||
let range = c.req.header('range');
|
||||
const ifRange = c.req.header('if-range');
|
||||
if (range && ifRange && etag && ifRange.trim() !== etag) {
|
||||
// Stale partial: send the whole current file. As a stream, because Bun
|
||||
// applies the request's Range to a Bun.file body on its own.
|
||||
return new Response(file.stream(), { status: 200, headers: { ...headers, 'Content-Length': String(total) } });
|
||||
}
|
||||
if (!range) return new Response(file, { status: 200, headers: { ...headers, 'Content-Length': String(total) } });
|
||||
const m = /^bytes=(\d*)-(\d*)$/.exec(range.trim());
|
||||
let start, end;
|
||||
if (m && m[1] !== '') {
|
||||
@@ -1524,23 +1535,73 @@ app.post('/api/media/:id/redownload', async (c) => {
|
||||
});
|
||||
|
||||
// Stream the server-cached copy as a save (OPFS writes it on the device).
|
||||
function cachedDownloadResponse(videoId, fp, row) {
|
||||
const file = Bun.file(`${MEDIA_DIR}/${videoId}.${row.gen}.mp4`);
|
||||
if (fp) recordVideoAccess(fp, { id: videoId }).catch(() => {});
|
||||
// Saves fetch the copy in byte ranges (docs/resumable-downloads-plan.md):
|
||||
// the ETag pins the generation, so a resumed save never mixes two copies.
|
||||
const prepareSkips = new Map(); // videoId → { reason, at } for /prepare
|
||||
const DOWNLOAD_EXPOSE = 'X-Content-SHA256, Content-Range, Content-Length, ETag, Accept-Ranges';
|
||||
function cachedDownloadResponse(c, videoId, fp, row) {
|
||||
if (fp && !c.req.header('range')) recordVideoAccess(fp, { id: videoId }).catch(() => {});
|
||||
media.touch(videoId);
|
||||
return new Response(file, {
|
||||
status: 200,
|
||||
return rangeFileResponse(c, `${MEDIA_DIR}/${videoId}.${row.gen}.mp4`, 'video/mp4', 'no-store', {
|
||||
etag: `"${videoId}.${row.gen}"`,
|
||||
headers: {
|
||||
'Content-Type': 'video/mp4',
|
||||
'Content-Length': String(file.size),
|
||||
'Content-Disposition': `attachment; filename="${videoId}.mp4"`,
|
||||
'Cache-Control': 'no-store',
|
||||
'Access-Control-Allow-Origin': '*',
|
||||
...(row.sha256 ? { 'X-Content-SHA256': row.sha256, 'Access-Control-Expose-Headers': 'X-Content-SHA256' } : {}),
|
||||
'Access-Control-Expose-Headers': DOWNLOAD_EXPOSE,
|
||||
...(row.sha256 ? { 'X-Content-SHA256': row.sha256 } : {}),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
// GET /api/download/:videoId/prepare[?hevc=1] — never blocks. Starts (or
|
||||
// joins) the server-side fetch and reports where it is:
|
||||
// ready → { gen, size, sha256 }: fetch /api/download/:id in ranges
|
||||
// working → poll again (the server is still getting it from YouTube)
|
||||
// legacy → this copy can't be served in ranges (too long for the cache,
|
||||
// HEVC the device can't play, cache offline): use the old
|
||||
// single-request save
|
||||
app.get('/api/download/:videoId/prepare', async (c) => {
|
||||
const videoId = (c.req.param('videoId') || '').trim();
|
||||
const nocache = { 'Cache-Control': 'no-store' };
|
||||
if (isUpload(videoId)) {
|
||||
const u = await notesDb.getUpload(videoId);
|
||||
if (!u) return c.json({ ok: false, error: 'upload not found' }, 404);
|
||||
let size = 0;
|
||||
try { size = Bun.file(uploads.filePath(u)).size; } catch { /* missing */ }
|
||||
return c.json({ ok: true, state: 'ready', gen: 0, size, sha256: null, etag: `"${u.id}"`, ext: u.ext }, 200, nocache);
|
||||
}
|
||||
if (!media.isMediaId(videoId)) return c.json({ ok: false, error: 'invalid video id' }, 400);
|
||||
const hevcOk = c.req.query('hevc') === '1';
|
||||
const ready = await media.getReady(videoId);
|
||||
if (ready) {
|
||||
if (!servableTo(ready, hevcOk)) return c.json({ ok: true, state: 'legacy', reason: 'hevc' }, 200, nocache);
|
||||
media.touch(videoId);
|
||||
return c.json({ ok: true, state: 'ready', gen: ready.gen, size: ready.size, sha256: ready.sha256 || null,
|
||||
etag: `"${videoId}.${ready.gen}"`, ext: 'mp4' }, 200, nocache);
|
||||
}
|
||||
// A job refused moments ago (too long for the cache, cache offline…) would
|
||||
// just be refused again — answer from memory instead of re-probing YouTube.
|
||||
const prior = prepareSkips.get(videoId);
|
||||
if (prior && Date.now() - prior.at < 10 * 60_000) return c.json({ ok: true, state: 'legacy', reason: prior.reason }, 200, nocache);
|
||||
// Start or join the job; its result is picked up by the next poll.
|
||||
let skip = null;
|
||||
media.ensureCached(videoId, { priority: HIGH }).catch((err) => {
|
||||
skip = err;
|
||||
if (err instanceof MediaSkip) {
|
||||
prepareSkips.set(videoId, { reason: err.message, at: Date.now() });
|
||||
if (prepareSkips.size > 2000) prepareSkips.clear();
|
||||
}
|
||||
});
|
||||
await new Promise((r) => setTimeout(r, 0));
|
||||
if (skip) {
|
||||
const legacy = skip instanceof MediaSkip;
|
||||
return c.json({ ok: legacy, state: legacy ? 'legacy' : 'failed', reason: skip.message }, legacy ? 200 : 500, nocache);
|
||||
}
|
||||
const st = await media.status(videoId);
|
||||
if (st.status === 'failed') return c.json({ ok: true, state: 'legacy', reason: st.error || 'server fetch failed' }, 200, nocache);
|
||||
return c.json({ ok: true, state: 'working', status: st.status }, 200, nocache);
|
||||
});
|
||||
|
||||
// GET /api/download/:videoId
|
||||
// Saves go through the server cache: the fetch runs as a server-owned job
|
||||
// (a disconnecting client no longer kills it — the copy lands for next
|
||||
@@ -1555,11 +1616,10 @@ app.get('/api/download/:videoId', async (c) => {
|
||||
if (isUpload(videoId)) {
|
||||
const u = await notesDb.getUpload(videoId);
|
||||
if (!u) return c.json({ ok: false, error: 'upload not found' }, 404);
|
||||
const f = Bun.file(uploads.filePath(u));
|
||||
return new Response(f, { status: 200, headers: {
|
||||
'Content-Type': u.mime, 'Content-Length': String(f.size),
|
||||
'Content-Disposition': `attachment; filename="${u.id}.${u.ext}"`, 'Cache-Control': 'no-store',
|
||||
} });
|
||||
return rangeFileResponse(c, uploads.filePath(u), u.mime, 'no-store', {
|
||||
etag: `"${u.id}"`,
|
||||
headers: { 'Content-Disposition': `attachment; filename="${u.id}.${u.ext}"`, 'Access-Control-Expose-Headers': DOWNLOAD_EXPOSE },
|
||||
});
|
||||
}
|
||||
const fp = c.req.query('fp');
|
||||
|
||||
@@ -1586,7 +1646,7 @@ app.get('/api/download/:videoId', async (c) => {
|
||||
const row = await media.ensureCached(videoId, { priority: HIGH });
|
||||
// A device that can't decode HEVC must not save the HEVC copy — give it
|
||||
// the legacy H.264 save instead (the ?mux=1 path included).
|
||||
if (servableTo(row, c.req.query('hevc') === '1')) return cachedDownloadResponse(videoId, fp, row);
|
||||
if (servableTo(row, c.req.query('hevc') === '1')) return cachedDownloadResponse(c, videoId, fp, row);
|
||||
cacheErr = Object.assign(new Error('cached copy is HEVC; client did not ask for it'), { code: 'SKIPPED' });
|
||||
} catch (err) {
|
||||
cacheErr = err;
|
||||
|
||||
Reference in New Issue
Block a user