Files
ytplayer/frontend/direct-media.js

169 lines
13 KiB
JavaScript

/* Paired-room transport; signalling only goes to the server. */
(function (root) {
'use strict';
const waiting = new Map();
let profilePeers = [];
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');
if (pending.has(`${key}:${to}:${meta.id}`) || [...transfers.values()].some(s => s.key === key && s.peer === to && s.claim.id === meta.id)) return;
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, extension: file.name.split('.').pop().toLowerCase(), streamMime: await root.DirectStream?.inspect(file), title: meta.title || meta.id });
if (!claim) throw Error('Unsupported file metadata');
const pendingKey = `${key}:${to}:${claim.id}`;
const item = { file, claim }; pending.set(pendingKey, item);
setTimeout(() => { if (pending.get(pendingKey) === item) pending.delete(pendingKey); }, 60000);
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 (error) s.preview?.close(); else s.preview?.end();
if (s.worker) { s.worker.postMessage({ op: 'abort' }); setTimeout(() => s.worker.terminate(), 1000); }
if (error) { waiting.get(s.claim.id)?.resolve(false); waiting.delete(s.claim.id); 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) {
clearTimeout(s.timer);
const pc = s.pc = new RTCPeerConnection({ iceServers: [{ urls: 'stun:stun.l.google.com:19302' }] });
s.ice = 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 && !s.verifying) setTimeout(() => { if (!s.done && !s.verifying) 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.verifying = true; clearTimeout(s.idle); s.worker.postMessage({ op: 'finish' }); } else throw Error('Unexpected file control'); }
else { s.preview?.push(e.data); 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' })); clearTimeout(s.idle); s.idle = setTimeout(() => finish(s, 'Receiver did not finish verification'), 30 * 60000); }
}
} 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);
s.preview = m.file.streamMime ? root.DirectStream?.create(m.file.streamMime, m.file.title) : null;
const wait = waiting.get(m.file.id); if (wait) clearTimeout(wait.timer);
send(key, { type: 'direct', action: 'accept', token: m.token });
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;
connection(s, true);
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); waiting.get(result.file.id)?.resolve(true); waiting.delete(result.file.id); publish(); 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;
if (m.action === 'request') { if (enabled() && hooks.settings().p2pShare !== false) offer(key, m.from, hooks.meta?.(m.id) || { id: m.id, title: m.id }).catch(e => hooks.notify(e.message)); 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}`);
const s = { key, token: m.token, peer: m.to, claim: p.claim, file: p.file }; transfers.set(m.token, s); s.timer = setTimeout(() => finish(s), 60000); return;
}
const s = transfers.get(m.token); if (!s || s.key !== key || m.from !== s.peer) return;
if (m.action === 'accept') { s.accepted = true; clearTimeout(s.timer); s.timer = setTimeout(() => finish(s, 'Receiver stopped preparing the file'), 30 * 60000); }
else if (m.action === 'decline' || m.action === 'complete') finish(s);
else if (m.action === 'signal') {
const d = m.data;
if (!s.pc && !s.receiver && s.accepted) { if (d.kind === 'ice') { (s.ice ||= []).push(d.candidate); return; } if (d.kind === 'offer') connection(s, false); }
if (!s.pc) return;
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 = () => { waiting.get(m.file.id)?.resolve(false); waiting.delete(m.file.id); 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();
}
async function publish() {
if (!root.P2PClient?.isConnected()) return;
const files = await root.OPFS.listVideos({ strict: true });
root.P2PClient.send({ type: 'direct-inventory', ids: enabled() && hooks.settings().p2pShare !== false ? files.map(f => f.id) : [] });
}
function source(id) { if (!enabled()) return null; for (const [key, room] of rooms) { const peer = room.peers().find(p => p.files?.includes(id)); if (peer) return { key, peer }; } return null; }
function hasSource(id) { return !!source(id); }
function obtain(id) {
if (!hasSource(id)) return Promise.resolve(null);
if (waiting.has(id)) return waiting.get(id).promise;
const { key, peer } = source(id);
let resolve; const promise = new Promise(r => { resolve = r; });
const wait = { promise, resolve }; waiting.set(id, wait);
send(key, { type: 'direct', action: 'request', to: peer.id, id });
wait.timer = setTimeout(() => { if (waiting.get(id)?.promise === promise) { waiting.delete(id); resolve(false); hooks.notify('The other device did not accept or is unavailable'); hooks.fallback(hooks.meta?.(id) || { id, title: id }); } }, 60000);
return promise;
}
function configure(h) {
hooks = { ...hooks, ...h };
if (root.P2PClient) {
root.P2PClient.onMessage('hello', () => { attach('profile', { send: root.P2PClient.send, peers: () => profilePeers }); publish(); });
root.P2PClient.onMessage('direct-peers', m => { profilePeers = m.peers || []; });
root.P2PClient.onMessage('direct', m => handle('profile', m));
}
root.SettingsSections.register({ id: 'direct-transfer', title: 'Direct device transfer', cluster: 'Library & storage', summary: () => enabled() ? 'On · paired devices only' : 'Off or unsupported', icon: '<path d="M4 6h6v12H4zM14 6h6v12h-6M10 10h4m-4 4h4"/>', 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(); publish(); }; 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, obtain, hasSource, publish };
}(globalThis));