From 08ff8398b579c9d84dea55acd7440a9abf323e96 Mon Sep 17 00:00:00 2001 From: Jonathan Sykes Date: Sat, 3 Oct 2026 20:32:45 +0800 Subject: [PATCH] Transfer saved songs directly between paired devices with resumable OPFS writes --- frontend/app.js | 10 +++ frontend/direct-media.css | 6 ++ frontend/direct-media.js | 134 ++++++++++++++++++++++++++++++++ frontend/direct-recv-worker.js | 38 +++++++++ frontend/index.html | 3 + frontend/p2p-client.js | 2 +- frontend/sw.js | 4 + playwright.direct.config.js | 2 + tests/direct-transfer.spec.js | 28 +++++++ tests/fixtures/direct-server.js | 27 +++++++ 10 files changed, 253 insertions(+), 1 deletion(-) create mode 100644 frontend/direct-media.css create mode 100644 frontend/direct-media.js create mode 100644 frontend/direct-recv-worker.js create mode 100644 playwright.direct.config.js create mode 100644 tests/direct-transfer.spec.js create mode 100644 tests/fixtures/direct-server.js diff --git a/frontend/app.js b/frontend/app.js index 6bf57ca..7e64f63 100755 --- a/frontend/app.js +++ b/frontend/app.js @@ -4754,6 +4754,8 @@ const Remote = (() => { const hostSend = (m) => { if (host.ws && host.ws.readyState === 1) host.ws.send(JSON.stringify(m)); }; function onHostMessage(m) { + DirectMedia.attach('remote-host', { send: hostSend, peers: () => host.remotes.map(r => ({ id: r.rid, name: r.name })) }); + if (m.type === 'direct') { DirectMedia.handle('remote-host', m); return; } if (m.type === 'hello') { host.code = m.code; host.codeExpires = m.codeExpires; host.remotes = m.remotes || []; host.lastState = ''; host.lastQueue = ''; @@ -5061,6 +5063,8 @@ const Remote = (() => { } function onRemoteMessage(m) { + DirectMedia.attach('remote-client', { send: rcSend, peers: () => [{ id: 'host', name: 'Paired screen' }] }); + if (m.type === 'direct') { DirectMedia.handle('remote-client', m); return; } if (m.type === 'hello') { rc.online = true; $('rvHostName').textContent = m.hostName || 'Screen'; @@ -7452,6 +7456,8 @@ const Party = (() => { } function onMessage(m) { + DirectMedia.attach('party', { send, peers: () => st.members.filter(m => m.pid !== pid).map(m => ({ id: m.pid, name: m.name })) }); + if (m.type === 'direct') { DirectMedia.handle('party', m); return; } if (m.type === 'hello') { st.code = m.code; if (m.secret) { st.secret = m.secret; ls.set(HOST_KEY, JSON.stringify({ code: m.code, secret: m.secret })); } @@ -12527,6 +12533,10 @@ async function boot() { Party.boot(); PartyDJ.register(() => $('navPartyBtn').click()); SetlistImport.configure({ library: async () => { let uploads = []; try { const j = await webFetch('/api/uploads?limit=200'); uploads = j.uploads || []; } catch {} return [...data.history, ...data.playlists.flatMap(p => p.videos || []), ...uploads]; }, search: async q => (await API.search(q)).results || [], modal: showModal, close: closeModal, create: (name, videos) => { const pl = { id: uid(), name, videos: videos.map(v => ({ ...slim(v), tags: v.tags, key: v.key, chordPro: v.chordPro })) }; data.playlists.push(pl); persist(); view = { type: 'playlist', id: pl.id }; render(); toast(`Imported ${videos.length} songs`); } }); + DirectMedia.configure({ settings: () => data.settings, persist, current: () => current?.meta, notify: toast, + saved: async file => { cachedIds.add(file.id); await DeviceDB.putFile({ videoId: file.id, cid: file.cid, size: file.size, savedAt: Date.now(), state: 'unverified' }); if (!data.history.some(v => v.id === file.id)) data.history.unshift({ ...file, thumbnail: '' }); persist(); renderSidebar(); }, + fallback: file => { const body = document.createElement('p'); body.textContent = 'The devices could not finish a direct connection. Keep both screens open and send again to resume, or choose the existing server download.'; showModal('Direct transfer unavailable', body, [{ label: 'Keep partial copy', onClick: closeModal }, { label: 'Use server copy', onClick: () => { closeModal(); preload(file, { retry: true, directFallback: true }); } }]); }, + }); Remote.boot(); LyricsDisplay.register(); LowerThird.register(() => $('remoteChip').click()); diff --git a/frontend/direct-media.css b/frontend/direct-media.css new file mode 100644 index 0000000..1b3742b --- /dev/null +++ b/frontend/direct-media.css @@ -0,0 +1,6 @@ +.direct-transfer-dialog{max-width:min(420px,calc(100vw - 32px));padding:24px;border:1px solid var(--border);border-radius:var(--radius,16px);background:var(--surface,#202026);color:var(--text,#fff)} +.direct-transfer-dialog::backdrop{background:rgba(0,0,0,.65)} +.direct-transfer-dialog p{line-height:1.6;color:var(--text-2);overflow-wrap:anywhere} +.direct-transfer-dialog button,#settings-section-direct-transfer button{min-height:44px;padding:10px 16px;margin:4px;border:1px solid var(--border);border-radius:10px;background:var(--surface-2,#303039);color:var(--text,#fff)} +.direct-transfer-dialog button:focus-visible,#settings-section-direct-transfer button:focus-visible{outline:2px solid var(--accent);outline-offset:3px} +#settings-section-direct-transfer label{display:flex;align-items:center;gap:12px;min-height:44px} diff --git a/frontend/direct-media.js b/frontend/direct-media.js new file mode 100644 index 0000000..7077115 --- /dev/null +++ b/frontend/direct-media.js @@ -0,0 +1,134 @@ +/* Paired-room transport; signalling only goes to the server. */ +(function (root) { + 'use strict'; + const rooms = new Map(), transfers = new Map(), pending = new Map(); + let hooks = { settings: () => ({}), persist() {}, current: () => null, saved() {}, notify() {}, fallback() {} }; + const enabled = () => hooks.settings().directTransfer !== false && !!root.RTCPeerConnection; + function attach(key, room) { rooms.set(key, room); } + function send(key, m) { const r = rooms.get(key); if (!r) throw Error('Paired device disconnected'); r.send(m); } + function signal(key, token, data) { send(key, { type: 'direct', action: 'signal', token, data }); } + async function offer(key, to, meta) { + if (!enabled()) throw Error('Direct device transfer is unavailable or turned off'); + if (!meta?.id) throw Error('Choose a saved song first'); + const file = await root.OPFS.getFileObject(meta.id); + if (!file) throw Error('Save this song on this device before sending it'); + const hash = root.Sha256.create(); + for (let at = 0; at < file.size; at += 1024 * 1024) hash.update(new Uint8Array(await file.slice(at, at + 1024 * 1024).arrayBuffer())); + const claim = root.DirectProtocol.file({ ...meta, cid: hash.hex(), size: file.size, title: meta.title || meta.id }); + if (!claim) throw Error('Unsupported file metadata'); + pending.set(`${key}:${to}:${claim.id}`, { file, claim }); + send(key, { type: 'direct', action: 'invite', to, file: claim }); + hooks.notify('Waiting for the receiving device to accept'); + } + function finish(s, error) { + if (s.done) return; s.done = true; clearTimeout(s.timer); clearTimeout(s.idle); + s.pc?.close(); transfers.delete(s.token); + if (s.worker) { s.worker.postMessage({ op: 'abort' }); setTimeout(() => s.worker.terminate(), 1000); } + if (error) { hooks.notify(`${error}. Partial copy kept for resume.`); hooks.fallback(s.claim, () => offer(s.key, s.peer, s.claim).catch(e => hooks.notify(e.message))); } + } + function connection(s, receiver) { + const pc = s.pc = new RTCPeerConnection({ iceServers: [{ urls: 'stun:stun.l.google.com:19302' }] }); + s.ice = []; + pc.onicecandidate = e => { if (e.candidate) signal(s.key, s.token, { kind: 'ice', candidate: e.candidate.toJSON() }); }; + s.timer = setTimeout(() => finish(s, 'Direct connection failed after 10 seconds'), 10000); + pc.onconnectionstatechange = () => { if (['failed', 'closed'].includes(pc.connectionState)) finish(s, 'Direct connection closed'); }; + const bump = () => { clearTimeout(s.idle); s.idle = setTimeout(() => finish(s, 'Direct transfer stalled'), 30000); }; + const channel = dc => { + s.dc = dc; dc.binaryType = 'arraybuffer'; dc.bufferedAmountLowThreshold = 1024 * 1024; + dc.onopen = () => { clearTimeout(s.timer); bump(); if (receiver && s.offset != null) dc.send(JSON.stringify({ t: 'get', offset: s.offset })); }; + dc.onclose = () => { if (!s.done) setTimeout(() => { if (!s.done) finish(s, 'Direct channel closed'); }, 3000); }; + dc.onmessage = async e => { + bump(); + try { + if (receiver) { + if (typeof e.data === 'string') { const m = JSON.parse(e.data); if (m.t === 'end') s.worker.postMessage({ op: 'finish' }); else throw Error('Unexpected file control'); } + else s.worker.postMessage({ op: 'chunk', buf: e.data }, [e.data]); + } else { + if (typeof e.data !== 'string' || s.started) throw Error('Unexpected file request'); + const m = JSON.parse(e.data); if (m.t !== 'get') throw Error('Unexpected file request'); + let pos = root.DirectProtocol.offset(m.offset, s.file.size); s.started = true; + while (pos < s.file.size && !s.done) { + if (dc.bufferedAmount > 4 * 1024 * 1024) await new Promise((resolve, reject) => { + const timeout = setTimeout(() => { cleanup(); reject(Error('Peer stopped receiving')); }, 10000); + const cleanup = () => { clearTimeout(timeout); dc.removeEventListener('bufferedamountlow', low); dc.removeEventListener('close', closed); }; + const low = () => { cleanup(); resolve(); }, closed = () => { cleanup(); reject(Error('Peer disconnected')); }; + dc.addEventListener('bufferedamountlow', low); dc.addEventListener('close', closed); + }); + const range = root.DirectProtocol.nextChunk(pos, s.file.size); + dc.send(await s.file.slice(range.start, range.end).arrayBuffer()); pos = range.end; bump(); + } + if (!s.done) dc.send(JSON.stringify({ t: 'end' })); + } + } catch (err) { finish(s, err.message); } + }; + }; + if (receiver) channel(pc.createDataChannel('file', { ordered: true })); else pc.ondatachannel = e => channel(e.channel); + return pc; + } + async function accept(key, m) { + const s = { key, token: m.token, claim: m.file, peer: m.from, receiver: true }; transfers.set(m.token, s); + connection(s, true); + s.worker = new Worker('/direct-recv-worker.js'); + s.worker.onmessage = async e => { + const result = e.data; + if (result.op === 'opened') { + s.offset = result.offset; + send(key, { type: 'direct', action: 'accept', token: m.token }); + const offer = await s.pc.createOffer(); await s.pc.setLocalDescription(offer); signal(key, m.token, { kind: 'offer', sdp: offer.sdp }); + } else if (result.op === 'done') { + send(key, { type: 'direct', action: 'complete', token: m.token }); + await hooks.saved(result.file); hooks.notify(`Saved ${result.file.title} directly from the other device`); finish(s); + } else if (result.op === 'error') finish(s, result.error); + }; + s.worker.postMessage({ op: 'open', file: m.file }); + } + async function handle(key, m) { + if (m.type !== 'direct') return; + try { + if (m.action === 'invite') { + if (!enabled() || !root.DirectProtocol.file(m.file)) return send(key, { type: 'direct', action: 'decline', token: m.token }); + confirmReceive(key, m); return; + } + if (m.action === 'invited') { + const p = pending.get(`${key}:${m.to}:${m.file.id}`); if (!p) return; + pending.delete(`${key}:${m.to}:${m.file.id}`); + transfers.set(m.token, { key, token: m.token, peer: m.to, claim: p.claim, file: p.file }); return; + } + const s = transfers.get(m.token); if (!s || s.key !== key || m.from !== s.peer) return; + if (m.action === 'accept') connection(s, false); + else if (m.action === 'decline' || m.action === 'complete') finish(s); + else if (m.action === 'signal' && s.pc) { + const d = m.data; + if (d.kind === 'ice') { if (s.pc.remoteDescription) await s.pc.addIceCandidate(d.candidate); else s.ice.push(d.candidate); } + else { + await s.pc.setRemoteDescription({ type: d.kind, sdp: d.sdp }); + for (const c of s.ice.splice(0)) await s.pc.addIceCandidate(c); + if (d.kind === 'offer') { const a = await s.pc.createAnswer(); await s.pc.setLocalDescription(a); signal(key, m.token, { kind: 'answer', sdp: a.sdp }); } + } + } + } catch (e) { const s = transfers.get(m.token); if (s) finish(s, e.message); else hooks.notify(e.message); } + } + function confirmReceive(key, m) { + const dialog = document.createElement('dialog'); dialog.className = 'direct-transfer-dialog'; + const heading = document.createElement('h2'); heading.textContent = 'Save from your paired device?'; + const text = document.createElement('p'); text.textContent = `${m.file.title} · ${(m.file.size / 1024 ** 2).toFixed(1)} MB. Media travels directly between devices. Keep both screens open. Playback starts after the full copy is verified.`; + const yes = document.createElement('button'); yes.textContent = 'Accept & save'; + const no = document.createElement('button'); no.textContent = 'Decline'; + const decline = () => { send(key, { type: 'direct', action: 'decline', token: m.token }); dialog.remove(); }; + yes.onclick = () => { dialog.remove(); accept(key, m).catch(e => hooks.notify(e.message)); }; no.onclick = decline; dialog.oncancel = decline; + dialog.append(heading, text, yes, no); document.body.append(dialog); dialog.showModal(); + } + function configure(h) { + hooks = { ...hooks, ...h }; + root.SettingsSections.register({ id: 'direct-transfer', title: 'Direct device transfer', cluster: 'Library & storage', summary: () => enabled() ? 'On · paired devices only' : 'Off or unsupported', icon: '', render(container) { + const label = document.createElement('label'); const toggle = document.createElement('input'); toggle.type = 'checkbox'; toggle.checked = hooks.settings().directTransfer !== false; + toggle.onchange = () => { hooks.settings().directTransfer = toggle.checked; hooks.persist(); }; label.append(toggle, ' Direct device transfer (P2P)'); + const help = document.createElement('p'); help.textContent = 'Pair through Remote or join a Watch Party, then send your currently playing saved song. The receiver confirms before saving. Keep both devices open. No TURN relay is used.'; + container.append(label, help); + for (const [key, r] of rooms) for (const peer of r.peers()) { + const b = document.createElement('button'); b.textContent = `Send saved song to ${peer.name}`; b.style.minHeight = '44px'; b.onclick = () => offer(key, peer.id, hooks.current()).catch(e => hooks.notify(e.message)); container.append(b); + } + } }); + } + root.DirectMedia = { configure, attach, handle, offer, enabled }; +}(globalThis)); diff --git a/frontend/direct-recv-worker.js b/frontend/direct-recv-worker.js new file mode 100644 index 0000000..aeafe54 --- /dev/null +++ b/frontend/direct-recv-worker.js @@ -0,0 +1,38 @@ +'use strict'; +importScripts('/sha256.js'); +let dir, handle, access, hash, pos = 0, claim, name; +let chain = Promise.resolve(); +self.onmessage = e => { chain = chain.then(() => run(e.data)); }; +async function close() { if (access) { access.flush(); access.close(); access = null; } } +async function run(m) { + try { + if (m.op === 'open') { + claim = m.file; + if (!/^[A-Za-z0-9_-]{1,128}$/.test(claim.id) || !/^[a-f0-9]{64}$/.test(claim.cid)) throw Error('Invalid file identity'); + dir = await (await navigator.storage.getDirectory()).getDirectoryHandle('videos', { create: true }); + name = `${claim.id}.${claim.cid}.direct.part`; + handle = await dir.getFileHandle(name, { create: true }); + access = await handle.createSyncAccessHandle(); + pos = access.getSize(); if (pos > claim.size) { access.truncate(0); pos = 0; } + hash = self.Sha256.create(); + for (let at = 0; at < pos; at += 1024 * 1024) { const b = new Uint8Array(Math.min(1024 * 1024, pos - at)); access.read(b, { at }); hash.update(b); } + self.postMessage({ op: 'opened', offset: pos }); + } else if (m.op === 'chunk') { + const b = new Uint8Array(m.buf); + if (!access || b.length > 65536 || pos + b.length > claim.size) throw Error('Invalid file chunk'); + if (access.write(b, { at: pos }) !== b.length) throw Error('Incomplete OPFS write'); + hash.update(b); pos += b.length; access.flush(); + self.postMessage({ op: 'progress', received: pos }); + } else if (m.op === 'abort') { await close(); self.postMessage({ op: 'closed' }); } + else if (m.op === 'finish') { + await close(); + if (pos !== claim.size || hash.hex() !== claim.cid) { await dir.removeEntry(name); throw Error('File verification failed; retry starts from zero'); } + const final = `${claim.id}.mp4`; + // Copy only after verification; existing saved playback remains intact until then. + const out = await (await dir.getFileHandle(final, { create: true })).createSyncAccessHandle(); + try { const f = await handle.getFile(); for (let at = 0; at < f.size; at += 1024 * 1024) out.write(new Uint8Array(await f.slice(at, at + 1024 * 1024).arrayBuffer()), { at }); out.truncate(f.size); out.flush(); } finally { out.close(); } + await dir.removeEntry(name); + self.postMessage({ op: 'done', file: claim }); + } + } catch (e) { await close(); self.postMessage({ op: 'error', error: e.message }); } +} diff --git a/frontend/index.html b/frontend/index.html index faee07f..67fd73f 100755 --- a/frontend/index.html +++ b/frontend/index.html @@ -28,6 +28,7 @@ + @@ -660,6 +661,8 @@ + + diff --git a/frontend/p2p-client.js b/frontend/p2p-client.js index 2e8a892..09e9d54 100644 --- a/frontend/p2p-client.js +++ b/frontend/p2p-client.js @@ -285,7 +285,7 @@ // and only for the exact cid this device holds. let restoring = false; onMessage('upload-request', async (m) => { - if (restoring || hooks.getSettings().p2pShare === false) return; + if (restoring || hooks.getSettings().p2pShare === false || hooks.getSettings().directTransfer !== false) return; const rec = window.DeviceDB ? await window.DeviceDB.getFile(m.videoId) : null; if (!rec || rec.cid !== m.cid) return; restoring = true; diff --git a/frontend/sw.js b/frontend/sw.js index 01e2297..5b9ee6f 100644 --- a/frontend/sw.js +++ b/frontend/sw.js @@ -81,6 +81,10 @@ const SHELL = [ '/hash-worker.js', '/p2p-client.js', '/p2p-core.js', + '/direct-media.css', + '/direct-protocol.js', + '/direct-media.js', + '/direct-recv-worker.js', '/p2p-transfer.js', '/p2p-recv-worker.js', '/settings-sections.js', diff --git a/playwright.direct.config.js b/playwright.direct.config.js new file mode 100644 index 0000000..ac16b52 --- /dev/null +++ b/playwright.direct.config.js @@ -0,0 +1,2 @@ +const { defineConfig } = require('@playwright/test'); +module.exports = defineConfig({ testDir: './tests', testMatch: /direct-transfer\.spec\.js/, timeout: 40000, workers: 1, use: { browserName: 'chromium', baseURL: 'http://localhost:8095', serviceWorkers: 'block' }, webServer: { command: 'bun tests/fixtures/direct-server.js', port: 8095, reuseExistingServer: false } }); diff --git a/tests/direct-transfer.spec.js b/tests/direct-transfer.spec.js new file mode 100644 index 0000000..4a47991 --- /dev/null +++ b/tests/direct-transfer.spec.js @@ -0,0 +1,28 @@ +const { test, expect } = require('@playwright/test'); +test('paired pages save verified bytes directly over WebRTC with metadata-only server logs', async ({ browser, request }) => { + const a = await browser.newContext(), b = await browser.newContext(); + const sender = await a.newPage(), receiver = await b.newPage(); + await sender.goto('/?id=sender'); await receiver.goto('/?id=receiver'); + await sender.waitForFunction(() => window.ready); await receiver.waitForFunction(() => window.ready); + await receiver.evaluate(async () => { + const bytes = Uint8Array.from({ length: 2000000 }, (_, i) => i % 251), h = Sha256.create(); h.update(bytes); + const dir = await (await navigator.storage.getDirectory()).getDirectoryHandle('videos', { create: true }); + const w = await (await dir.getFileHandle(`testvideo.${h.hex()}.direct.part`, { create: true })).createWritable(); + await w.write(bytes.slice(0, 12345)); await w.close(); + }); + await sender.locator('#send').click(); + await expect(receiver.getByRole('dialog')).toBeVisible(); + expect(await receiver.evaluate(() => !!window.received)).toBe(false); + await receiver.getByRole('button', { name: 'Accept & save' }).click(); + await receiver.waitForFunction(() => !!window.received); + const result = await receiver.evaluate(async () => { + const d = await (await navigator.storage.getDirectory()).getDirectoryHandle('videos'); + const f = await (await d.getFileHandle('testvideo.mp4')).getFile(); + const hash = Sha256.create(); hash.update(new Uint8Array(await f.arrayBuffer())); return { size: f.size, cid: hash.hex(), expected: window.received.cid }; + }); + expect(await receiver.evaluate(() => window.resumeOffsets)).toEqual([12345]); + expect(result.size).toBe(2000000); expect(result.cid).toBe(result.expected); + const logs = await (await request.get('/audit')).json(); + expect(logs.every(x => !x.http && ['invite','accept','signal','complete'].includes(x.signalling))).toBe(true); + await a.close(); await b.close(); +}); diff --git a/tests/fixtures/direct-server.js b/tests/fixtures/direct-server.js new file mode 100644 index 0000000..f51649e --- /dev/null +++ b/tests/fixtures/direct-server.js @@ -0,0 +1,27 @@ +import { createDirectRelay } from '../../server/direct-relay.js'; +const relay = createDirectRelay(), peers = new Map(), log = []; +Bun.serve({ port: 8095, + async fetch(req, server) { + const url = new URL(req.url); + if (url.pathname === '/ws') { const id = url.searchParams.get('id'); if (!['sender', 'receiver'].includes(id)) return new Response('test member required', { status: 403 }); server.upgrade(req, { data: { id } }); return; } + if (url.pathname === '/audit') return Response.json(log); + if (url.pathname.startsWith('/api/')) { log.push({ http: url.pathname }); return new Response('Media forbidden in direct fixture', { status: 500 }); } + if (url.pathname === '/') return new Response(`
`, { headers: { 'Content-Type': 'text/html' } }); + const path = url.pathname.slice(1); if (!/^[a-z0-9-]+\.js$/.test(path)) return new Response('not found', { status: 404 }); + return new Response(Bun.file(`frontend/${path}`)); + }, websocket: { + open(ws) { peers.set(ws.data.id, ws); }, + message(ws, message) { log.push({ signalling: typeof message === 'string' ? JSON.parse(message).action : 'BINARY' }); relay.handle(ws.data.id, typeof message === 'string' ? message : null, peers, (target,m) => target.send(JSON.stringify(m))); }, + close(ws) { peers.delete(ws.data.id); }, + } +});