Keep resume preparation and final verification outside connection timeouts
This commit is contained in:
@@ -35,7 +35,7 @@
|
||||
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 = 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'); };
|
||||
@@ -43,12 +43,12 @@
|
||||
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.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.worker.postMessage({ op: 'finish' }); else throw Error('Unexpected file control'); }
|
||||
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');
|
||||
@@ -64,7 +64,7 @@
|
||||
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' }));
|
||||
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); }
|
||||
};
|
||||
@@ -75,13 +75,14 @@
|
||||
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;
|
||||
connection(s, true);
|
||||
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;
|
||||
send(key, { type: 'direct', action: 'accept', token: m.token });
|
||||
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 });
|
||||
@@ -104,10 +105,12 @@
|
||||
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);
|
||||
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' && s.pc) {
|
||||
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 });
|
||||
@@ -139,9 +142,9 @@
|
||||
if (waiting.has(id)) return waiting.get(id).promise;
|
||||
const { key, peer } = source(id);
|
||||
let resolve; const promise = new Promise(r => { resolve = r; });
|
||||
waiting.set(id, { promise, resolve });
|
||||
const wait = { promise, resolve }; waiting.set(id, wait);
|
||||
send(key, { type: 'direct', action: 'request', to: peer.id, id });
|
||||
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);
|
||||
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) {
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
const { test, expect } = require('@playwright/test');
|
||||
test('paired pages save verified bytes directly over WebRTC with metadata-only server logs', async ({ browser, request }) => {
|
||||
for (const prepareDelay of [0, 12000]) test(`prepare delay ${prepareDelay}: 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.goto('/?id=sender'); await receiver.goto('/?id=receiver&prepareDelay=' + prepareDelay);
|
||||
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);
|
||||
|
||||
2
tests/fixtures/direct-server.js
vendored
2
tests/fixtures/direct-server.js
vendored
@@ -9,7 +9,7 @@ Bun.serve({ port: 8095,
|
||||
if (url.pathname === '/') return new Response(`<!doctype html><button id="send">Send</button><div id="status"></div><script src="/sha256-wasm.js"></script><script src="/sha256.js"></script><script src="/direct-protocol.js"></script><script src="/settings-sections.js"></script><script src="/direct-media.js"></script><script>
|
||||
const id = new URL(location.href).searchParams.get('id');
|
||||
const ws = new WebSocket('ws://localhost:8095/ws?id=' + id);
|
||||
window.resumeOffsets = []; const NativeWorker = window.Worker; window.Worker = class extends NativeWorker { constructor(...args) { super(...args); this.addEventListener('message', e => { if(e.data.op === 'opened') window.resumeOffsets.push(e.data.offset); }); } };
|
||||
window.resumeOffsets = []; const NativeWorker = window.Worker; window.Worker = class extends NativeWorker { constructor(...args) { super(...args); this.addEventListener('message', e => { if(e.data.op === 'opened') window.resumeOffsets.push(e.data.offset); }); } set onmessage(fn) { super.onmessage = e => { const delay = Number(new URL(location.href).searchParams.get('prepareDelay')) || 0; if(e.data.op === 'opened' && delay) setTimeout(() => fn(e), delay); else fn(e); }; } };
|
||||
const meta = { id: 'testvideo', title: 'Direct test', channel: 'Fixture', duration: 1 };
|
||||
window.OPFS = { getFileObject: async () => new File([Uint8Array.from({length: 2000000}, (_,i)=>i%251)], 'test.mp4') };
|
||||
DirectMedia.configure({ current: () => meta, notify: text => document.getElementById('status').textContent = text, saved: file => window.received = file, fallback: () => window.failed = true });
|
||||
|
||||
Reference in New Issue
Block a user