Close unconfirmed transfer paths and reject binary signalling envelopes

This commit is contained in:
Jonathan Sykes
2026-10-03 20:45:04 +08:00
parent 5661eea6e6
commit d78b06ff4b
6 changed files with 24 additions and 14 deletions

View File

@@ -3115,6 +3115,7 @@ async function renderAvailability(id) {
// Fetch a verified copy from another device (docs/p2p-architecture.md flow 6).
// Download-then-play: on success the file is saved offline like any save.
async function getFromPeers(videoObj, c, peers, onProgress) {
if (window.DirectMedia && !DirectMedia.enabled()) { toast('Direct device transfer is off. Enable it in Settings or use a server copy.'); return false; }
if (window.DirectMedia?.enabled()) { const ok = await DirectMedia.obtain(videoObj.id); if (!ok) toast('Pair the source device through Remote or use the same profile to receive its saved copy.'); return !!ok; }
const id = videoObj.id;
downloading.add(id);
@@ -12601,7 +12602,7 @@ async function boot() {
// Peer-to-peer: report what this device holds (on by default; settings.p2pShare).
if (WEB && window.P2PClient) {
window.P2PClient.start({ getSettings: () => data.settings, getProfile: () => (data.profile && data.profile.name) || '', getProfileSecret: () => ProfileSecret.getFor(data.profile?.name) });
if (window.P2PTransfer) window.P2PTransfer.start({ canShare: () => data.settings.p2pShare !== false });
if (window.P2PTransfer) window.P2PTransfer.start({ canShare: () => false }); // New scoped invitations replace the unconfirmed legacy sender.
}
data.playlists.forEach(preloadPlaylist);
preloadPinnedPlaylists();

View File

@@ -12,13 +12,16 @@
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');
pending.set(`${key}:${to}:${claim.id}`, { file, claim });
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');
}
@@ -30,6 +33,7 @@
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 = [];
pc.onicecandidate = e => { if (e.candidate) signal(s.key, s.token, { kind: 'ice', candidate: e.candidate.toJSON() }); };
@@ -88,7 +92,7 @@
}
async function handle(key, m) {
if (m.type !== 'direct') return;
if (m.action === 'request') { if (enabled()) offer(key, m.from, hooks.meta?.(m.id) || { id: m.id, title: m.id }).catch(e => hooks.notify(e.message)); 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 });
@@ -97,7 +101,7 @@
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 = { 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') connection(s, false);

View File

@@ -69,6 +69,7 @@ export function createP2pHub({ getDevice, authorizeProfile = async () => '', ena
}
async function message(ws, raw) {
if (typeof raw !== 'string') return;
const t = now();
ws.data.budget = ws.data.budget.filter((x) => t - x < 60_000);
if (ws.data.budget.length >= BUDGET_PER_MIN) { send(ws, { type: 'error', error: 'slow down' }); return; }
@@ -88,11 +89,7 @@ export function createP2pHub({ getDevice, authorizeProfile = async () => '', ena
if (!directRooms.has(profile)) directRooms.set(profile, createDirectRelay());
directRooms.get(profile).handle(ws.data.peer, m, peers, send); return;
}
if (m && m.type === 'signal' && typeof m.to === 'string') {
const target = online.get(byPeer.get(m.to));
if (!target || target === ws) { send(ws, { type: 'error', error: 'peer offline', to: m.to }); return; }
send(target, { type: 'signal', from: ws.data.peer, data: m.data });
}
if (m?.type === 'signal') send(ws, { type: 'error', error: 'use a confirmed direct-transfer invitation', to: m.to });
}
function close(ws) {

View File

@@ -51,7 +51,7 @@ test('auth on open, hello with opaque peer id, bad secret is closed', async () =
expect(hub.isOnline('dev_000000000000000b')).toBe(false);
});
test('signals are relayed between online peers only; close drops presence', async () => {
test('legacy unconfirmed signalling is refused; close drops presence', async () => {
const hub = createP2pHub({ getDevice: p2pDb.getDevice });
const a = fakeWs({ budget: [] });
const b = fakeWs({ budget: [] });
@@ -60,11 +60,12 @@ test('signals are relayed between online peers only; close drops presence', asyn
await hub.websocket.message(a, JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }));
await hub.websocket.message(b, JSON.stringify({ type: 'auth', device: 'dev_000000000000000b', secret: SECRET_B }));
await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: { sdp: 'x' } }));
expect(b.sent.at(-1)).toEqual({ type: 'signal', from: a.data.peer, data: { sdp: 'x' } });
expect(a.sent.at(-1)).toMatchObject({ type: 'error', error: 'use a confirmed direct-transfer invitation' });
expect(b.sent.some(m => m.type === 'signal')).toBe(false);
hub.websocket.close(b);
expect(hub.isOnline('dev_000000000000000b')).toBe(false);
await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: {} }));
expect(a.sent.at(-1)).toMatchObject({ type: 'error', error: 'peer offline' });
expect(a.sent.at(-1)).toMatchObject({ type: 'error', error: 'use a confirmed direct-transfer invitation' });
expect(hub.send('dev_000000000000000a', { type: 'ping' })).toBe(true);
expect(hub.send('dev_000000000000000b', { type: 'ping' })).toBe(false);
});
@@ -119,3 +120,10 @@ test('direct directory and invitations are restricted to authorized profile sock
await hub.websocket.message(b, JSON.stringify({ type: 'direct', action: 'request', to: peerIdOf('dev_000000000000000a'), id: 'localSong' }));
expect(a.sent.at(-1).action).toBe('request');
});
test('binary websocket payloads cannot serve as signalling envelopes', async () => {
const hub = createP2pHub({ getDevice: p2pDb.getDevice });
const ws = fakeWs({ budget: [] });
await hub.websocket.message(ws, Buffer.from(JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A })));
expect(ws.sent).toEqual([]); expect(ws.data.authed).not.toBe(true);
});

View File

@@ -170,7 +170,7 @@ export function createPartyHub({ now = () => Date.now() } = {}) {
return {
upgrade,
websocket: { maxPayloadLength: MAX_MSG, idleTimeout: 120, open, message: (ws, msg) => message(ws, typeof msg === 'string' ? msg : Buffer.from(msg).toString('utf8')), close },
websocket: { maxPayloadLength: MAX_MSG, idleTimeout: 120, open, message: (ws, msg) => message(ws, typeof msg === 'string' ? msg : null), close },
_parties: parties,
stop() { clearInterval(sweeper); },
};

View File

@@ -258,7 +258,7 @@ export function createRemoteHub({ requireSameNetwork = false } = {}) {
maxPayloadLength: MAX_MSG,
idleTimeout: 120,
open(ws) { if (ws.data.role === 'host') openHost(ws); else openRemote(ws); },
message(ws, msg) { onMessage(ws, typeof msg === 'string' ? msg : Buffer.from(msg).toString('utf8')); },
message(ws, msg) { onMessage(ws, typeof msg === 'string' ? msg : null); },
close(ws) { onClose(ws); },
};