From 75e437013fe570a514b784db02cbecb592618508 Mon Sep 17 00:00:00 2001 From: velzie Date: Mon, 21 Sep 2026 18:50:01 -0400 Subject: [PATCH] proper reconnect --- .../src/modules/Peer/PuterPeerConnection.js | 111 +++++++++++--- .../modules/Peer/PuterPeerConnection.test.js | 84 +++++++--- .../src/modules/Peer/PuterPeerServer.js | 135 +++++++++++++--- .../src/modules/Peer/PuterPeerServer.test.js | 144 +++++++++++++++++- src/puter-js/src/modules/Peer/events.js | 15 ++ src/puter-js/src/modules/Peer/signalling.js | 33 +++- 6 files changed, 448 insertions(+), 74 deletions(-) diff --git a/src/puter-js/src/modules/Peer/PuterPeerConnection.js b/src/puter-js/src/modules/Peer/PuterPeerConnection.js index c305674e9..48dc77588 100644 --- a/src/puter-js/src/modules/Peer/PuterPeerConnection.js +++ b/src/puter-js/src/modules/Peer/PuterPeerConnection.js @@ -13,8 +13,19 @@ import { /** @typedef {import('./types.js').PuterPeerOptions} PuterPeerOptions */ /** @typedef {import('./tracks.js').PuterPeerPublishOptions} PuterPeerPublishOptions */ -/** How many times ICE may be restarted before the connection is given up on. */ -const ICE_RESTART_LIMIT = 3; +/** + * How long to keep trying to rescue a link before giving it up. + * + * Counted in time rather than attempts, because what recovery is usually + * waiting on is the peer coming back - a lid reopened, a network rejoined - + * and that takes as long as it takes. A minute covers the ordinary cases + * and still ends a call to someone who is not coming back. + */ +const RECOVERY_BUDGET_MS = 60_000; + +/** Gap between restart attempts, doubling, so a long wait stays quiet. */ +const RECOVERY_BACKOFF_MS = 2000; +const RECOVERY_BACKOFF_MAX_MS = 15_000; /** How long an answered restart is given to bring the transport back. */ const TRANSPORT_SETTLE_TIMEOUT = 8000; @@ -52,10 +63,11 @@ export class PuterPeerConnection extends EventTarget { #tracks; #datachannel; #bufferedMessages = []; - #iceRestarts = 0; #recovering = false; - #peerGone = false; + /** The peer said goodbye, or closed the channel. Proof, not evidence. */ + #peerHungUp = false; #transportWaiters = new Set(); + #signallingWaiters = new Set(); /** * @param {object} peerConfig @@ -85,7 +97,7 @@ export class PuterPeerConnection extends EventTarget { }; this.#datachannel.onclose = () => { // A channel the peer closed cleanly is a hangup, not a fault. - this.#peerGone = true; + this.#peerHungUp = true; this.#doclose(undefined, undefined); }; this.#datachannel.onerror = (evt) => { @@ -107,8 +119,9 @@ export class PuterPeerConnection extends EventTarget { this.#channel.onanswer = (description, names) => this.#negotiator.acceptAnswer(description, names); this.#channel.oncandidate = (candidate) => this.#negotiator.acceptCandidate(candidate); this.#channel.onbye = (reason) => this.#onBye(reason); - this.#channel.onpeergone = (reason) => this.#onPeerGone(reason); + this.#channel.onpeergone = (reason, resumable) => this.#onPeerGone(reason, resumable); this.#channel.onunusable = () => this.#onSignallingLost(); + this.#channel.onusable = () => this.#wakeSignallingWaiters(); } /** @@ -141,19 +154,24 @@ export class PuterPeerConnection extends EventTarget { /** The peer said goodbye, so there is nothing to recover. */ #onBye ( reason ) { - this.#peerGone = true; + this.#peerHungUp = true; this.#doclose(reason, undefined); } /** - * The peer's signalling session ended. That is evidence the peer hung up, - * not proof - its data channel may still be carrying traffic - so it only - * decides anything once ICE has actually failed. + * The peer's signalling session ended. Evidence they hung up, never + * proof - the data channel may still be carrying traffic - so while we + * are connected it decides nothing and recovery is left to find out. + * A session the signaller says is reclaimable is not even evidence: + * the peer is expected back on it. * * @param {string} [reason] + * @param {boolean} [resumable] */ - #onPeerGone ( reason ) { - this.#peerGone = true; + #onPeerGone ( reason, resumable ) { + if ( ! resumable ) this.#peerHungUp = true; + // A handshake has nothing to wait on either way: whoever is dialling + // needs a definite answer rather than a socket that may come back. if ( ! this.connected ) this.#doclose(reason, undefined); } @@ -162,6 +180,9 @@ export class PuterPeerConnection extends EventTarget { * it returns, but an open data channel keeps working without it. */ #onSignallingLost () { + // Anything waiting on signalling has to look again: it may have + // gone for good rather than merely gone away. + this.#wakeSignallingWaiters(); if ( ! this.connected ) { this.#doclose('lost the signalling connection before the peer connected', undefined); } @@ -171,7 +192,6 @@ export class PuterPeerConnection extends EventTarget { this.#wakeTransportWaiters(); switch ( this.peerconnection.connectionState ) { case 'connected': - this.#iceRestarts = 0; this.#setLinkState('connected'); break; // 'disconnected' is transient: ICE either recovers by itself or @@ -216,16 +236,22 @@ export class PuterPeerConnection extends EventTarget { async #recover () { if ( this.closed || this.#recovering ) return; this.#recovering = true; + const deadline = Date.now() + RECOVERY_BUDGET_MS; + let attempt = 0; try { let unanswered = false; while ( ! this.closed && this.peerconnection.connectionState === 'failed' ) { - // Signalling does not come back on its own, so a restart - // that cannot be offered never will be. - if ( this.#peerGone || ! this.#channel.alive ) { + if ( this.#peerHungUp ) { + this.#doclose('the peer hung up', undefined); + return; + } + // No route left to offer over, and none coming. + if ( this.#channel.stranded ) { this.#doclose('the peer is no longer reachable', undefined); return; } - if ( this.#iceRestarts >= ICE_RESTART_LIMIT ) { + const left = deadline - Date.now(); + if ( left <= 0 ) { this.#doclose( unanswered ? 'the peer stopped responding' : 'could not restore the connection', undefined, @@ -233,19 +259,33 @@ export class PuterPeerConnection extends EventTarget { return; } - this.#iceRestarts++; - this.#setLinkState('recovering', { attempt: this.#iceRestarts, of: ICE_RESTART_LIMIT }); + // Nothing can be offered without signalling. It comes back + // on its own when a peer server reclaims its session, so + // wait on it rather than spend an attempt failing. + if ( ! this.#channel.alive ) { + await this.#awaitSignalling(left); + if ( ! this.#channel.alive ) continue; + } + + attempt++; + this.#setLinkState('recovering', { attempt }); try { await this.#negotiator.restartIce(); unanswered = false; } catch { unanswered = true; - continue; } - // Answered: the new candidates need a moment to either take - // over or fail, and only then is another attempt warranted. + + // Either the new candidates take over or they do not, and + // only then is another attempt worth making. if ( this.peerconnection.connectionState === 'failed' ) { - await this.#nextTransportChange(TRANSPORT_SETTLE_TIMEOUT); + const backoff = Math.min( + RECOVERY_BACKOFF_MAX_MS, + RECOVERY_BACKOFF_MS * 2 ** (attempt - 1), + ); + await this.#nextTransportChange( + Math.min(unanswered ? backoff : TRANSPORT_SETTLE_TIMEOUT, deadline - Date.now()), + ); } } } finally { @@ -253,6 +293,28 @@ export class PuterPeerConnection extends EventTarget { } } + /** + * Resolves when signalling is usable again, or when `timeout` runs out. + * + * @param {number} timeout + * @returns {Promise} + */ + #awaitSignalling ( timeout ) { + return new Promise((resolve) => { + const done = () => { + clearTimeout(timer); + this.#signallingWaiters.delete(done); + resolve(); + }; + const timer = setTimeout(done, timeout); + this.#signallingWaiters.add(done); + }); + } + + #wakeSignallingWaiters () { + for ( const wake of [...this.#signallingWaiters] ) wake(); + } + /** * Resolves on the next transport state change, or when `timeout` expires. * @@ -283,6 +345,7 @@ export class PuterPeerConnection extends EventTarget { // `close()` below need not raise another state change, so recovery is // told directly that there is nothing left to wait for. this.#wakeTransportWaiters(); + this.#wakeSignallingWaiters(); this.#setLinkState('closed'); this.#negotiator.stop(); this.#tracks.close(); @@ -290,7 +353,7 @@ export class PuterPeerConnection extends EventTarget { // Say goodbye while signalling is still up. An explicit hangup is the // only thing that separates a peer that left from one that broke, and // it carries the reason the far side reports to its own listeners. - if ( ! this.#peerGone ) this.#channel.sendBye(reason); + if ( ! this.#peerHungUp ) this.#channel.sendBye(reason); this.#channel.close(); if ( this.#datachannel ) { diff --git a/src/puter-js/src/modules/Peer/PuterPeerConnection.test.js b/src/puter-js/src/modules/Peer/PuterPeerConnection.test.js index f25f7185f..8f877dde7 100644 --- a/src/puter-js/src/modules/Peer/PuterPeerConnection.test.js +++ b/src/puter-js/src/modules/Peer/PuterPeerConnection.test.js @@ -82,19 +82,6 @@ describe('ICE recovery', () => { expect(pc.restarts).toBe(2); }); - it('gives up after the restart budget runs out', async () => { - const { conn, pc, channel, closes } = await makeConnection(); - - for ( let i = 0; i < 4; i++ ) { - pc.setConnectionState('failed'); - await answerOffer(channel); - } - - expect(pc.restarts).toBe(3); - expect(conn.closed).toBe(true); - expect(closes).toEqual(['could not restore the connection']); - }); - it('tries again when a restart goes unanswered', async () => { // A peer whose tab is throttled answers late, and the browser raises // no further state change while the transport stays failed, so one @@ -103,35 +90,82 @@ describe('ICE recovery', () => { const { conn, pc, closes } = await makeConnection(); pc.setConnectionState('failed'); - await vi.advanceTimersByTimeAsync(8000); + await vi.advanceTimersByTimeAsync(20_000); - expect(pc.restarts).toBe(2); + expect(pc.restarts).toBeGreaterThan(1); expect(conn.closed).toBe(false); expect(closes).toEqual([]); }); - it('treats a peer that answers no restart at all as gone', async () => { + it('keeps trying for a peer that could still come back', async () => { + // A lid closed for a minute is the ordinary case, and the side left + // awake must still be there when the other one wakes up. vi.useFakeTimers(); const { conn, pc, closes } = await makeConnection(); pc.setConnectionState('failed'); - await vi.advanceTimersByTimeAsync(8000 * 3); + await vi.advanceTimersByTimeAsync(45_000); + + expect(conn.closed).toBe(false); + expect(closes).toEqual([]); + }); + + it('gives the peer up once the budget is spent', async () => { + vi.useFakeTimers(); + const { conn, pc, closes } = await makeConnection(); + + pc.setConnectionState('failed'); + await vi.advanceTimersByTimeAsync(70_000); - expect(pc.restarts).toBe(3); expect(conn.closed).toBe(true); expect(closes).toEqual(['the peer stopped responding']); }); - it('does not bother restarting when the peer is known to be gone', async () => { + it('waits for signalling rather than spending attempts it cannot send', async () => { + vi.useFakeTimers(); const { conn, pc, channel, closes } = await makeConnection(); - channel.onpeergone('the peer went away'); + // The peer server's socket dropped; it is expected back on it. + channel.kill(); + pc.setConnectionState('failed'); + await vi.advanceTimersByTimeAsync(20_000); + + expect(pc.restarts).toBe(0); + expect(conn.closed).toBe(false); + expect(closes).toEqual([]); + + // It reclaimed its session, so the restart has somewhere to go. + channel.revive(); + channel.onusable(); + await vi.advanceTimersByTimeAsync(100); + + expect(pc.restarts).toBeGreaterThan(0); + expect(conn.closed).toBe(false); + }); + + it('does not bother restarting when the peer has hung up', async () => { + const { conn, pc, channel, closes } = await makeConnection(); + + channel.onpeergone('the peer went away', false); pc.setConnectionState('failed'); await flush(); expect(pc.restarts).toBe(0); expect(conn.closed).toBe(true); - expect(closes).toEqual(['the peer is no longer reachable']); + expect(closes).toEqual(['the peer hung up']); + }); + + it('keeps recovering when the peer only lost its signalling session', async () => { + vi.useFakeTimers(); + const { conn, pc, channel, closes } = await makeConnection(); + + channel.onpeergone('the peer server’s signalling dropped', true); + pc.setConnectionState('failed'); + await vi.advanceTimersByTimeAsync(10_000); + + expect(pc.restarts).toBeGreaterThan(0); + expect(conn.closed).toBe(false); + expect(closes).toEqual([]); }); }); @@ -214,10 +248,10 @@ describe('link state', () => { expect(closes).toEqual([]); }); - it('reports each restart attempt, and the budget it is spending', async () => { + it('reports each restart attempt', async () => { const { conn, pc, channel } = await makeConnection(); const seen = []; - conn.addEventListener('linkstate', (e) => seen.push([e.state, e.attempt, e.of])); + conn.addEventListener('linkstate', (e) => seen.push([e.state, e.attempt])); pc.setConnectionState('failed'); await answerOffer(channel); @@ -225,8 +259,8 @@ describe('link state', () => { await answerOffer(channel); expect(seen).toEqual([ - ['recovering', 1, 3], - ['recovering', 2, 3], + ['recovering', 1], + ['recovering', 2], ]); expect(conn.linkState).toBe('recovering'); }); diff --git a/src/puter-js/src/modules/Peer/PuterPeerServer.js b/src/puter-js/src/modules/Peer/PuterPeerServer.js index 6b60b9fbd..72aef6ffd 100644 --- a/src/puter-js/src/modules/Peer/PuterPeerServer.js +++ b/src/puter-js/src/modules/Peer/PuterPeerServer.js @@ -1,16 +1,33 @@ import { PuterPeerConnection } from './PuterPeerConnection.js'; -import { PuterPeerServerConnectionEvent } from './events.js'; +import { + PuterPeerServerConnectionEvent, + PuterPeerServerReconnectEvent, +} from './events.js'; import { ServerSignallingChannel } from './signalling.js'; /** @typedef {import('./types.js').PuterPeerOptions} PuterPeerOptions */ +/** How long to wait for the signaller to answer a registration. */ +const CREATE_TIMEOUT_MS = 15_000; + +/** Reconnect backoff for a socket that dropped under a running server. */ +const RECONNECT_BASE_MS = 500; +const RECONNECT_MAX_MS = 15_000; + /** * A peer server. One websocket to the signaller carries every client, so * connections are addressed on it by id; each gets its own channel wrapper. * - * The socket is registered once and never replaced. While it is down no new - * client can be accepted and nothing reaches the ones already connected, so - * `signallingAlive` says so and every channel is told. + * A socket that drops is dialled again, and the registration reclaims the + * session it left behind: same invite code, same clients still routed to + * it. That is what lets a host whose laptop slept pick the call back up + * instead of every guest having to build a new link. Until it is back, + * `signallingAlive` is false and every channel knows it. + * + * A reclaim that fails - gone too long, signaller restarted - registers + * fresh instead, under a new code. The clients that were attached to the + * old session can no longer be reached, so their channels are stranded + * rather than left to time out one signal at a time. */ export class PuterPeerServer extends EventTarget { connections = new Map(); @@ -24,6 +41,13 @@ export class PuterPeerServer extends EventTarget { /** @type {{ resolve: Function, reject: Function } | null} */ #pendingCreate = null; #alive = false; + #closed = false; + /** @type {PuterPeerOptions} */ + #options = {}; + /** @type {string | undefined} */ + #resumeToken; + #reconnectTimer = null; + #reconnectAttempts = 0; constructor ( peerConfig ) { super(); @@ -39,6 +63,19 @@ export class PuterPeerServer extends EventTarget { * @returns {Promise} */ async start ( options = {} ) { + this.#options = options; + const { inviteCode } = await this.#register(); + return inviteCode; + } + + /** + * Opens a socket and registers on it. On every call after the first this + * presents the resume token, so the signaller hands the same session + * back while it still holds it. + * + * @returns {Promise<{ inviteCode: string, resumed: boolean }>} + */ + async #register () { const ws = new WebSocket(this.#peerConfig.signallerUrl); this.#wsconn = ws; @@ -50,28 +87,36 @@ export class PuterPeerServer extends EventTarget { this.#alive = true; ws.onerror = null; - ws.onmessage = (event) => this.#message(event); + ws.onmessage = (event) => this.#message(event, ws); ws.onclose = () => { + if ( this.#wsconn !== ws ) return; this.#alive = false; this.#wsconn = null; for ( const channel of this.#channels.values() ) channel.onunusable(); - this.#settleCreate((pending) => { + this.#settleCreate(ws, (pending) => { pending.reject(new Error('Connection closed unexpectedly')); }); + if ( ! this.#closed ) this.#scheduleReconnect(); }; ws.send(JSON.stringify({ server: { create: { authToken: this.#peerConfig.authToken, - anonToken: options.anonToken, - port: options.port, + anonToken: this.#options.anonToken, + port: this.#options.port, + resume: this.#resumeToken, }, }, })); - this.inviteCode = await new Promise((resolve, reject) => { - this.#pendingCreate = { resolve, reject }; + const reply = await new Promise((resolve, reject) => { + const timer = setTimeout(() => { + this.#settleCreate(ws, (pending) => { + pending.reject(new Error('Server creation timed out')); + }); + }, CREATE_TIMEOUT_MS); + this.#pendingCreate = { ws, resolve, reject, timer }; }).catch((error) => { ws.onclose = null; try { @@ -79,18 +124,62 @@ export class PuterPeerServer extends EventTarget { } catch { // The failed registration is already unusable. } - this.#alive = false; - this.#wsconn = null; + // Only if this socket is still the live one: a registration + // already replaced must not mark its successor dead. + if ( this.#wsconn === ws ) { + this.#alive = false; + this.#wsconn = null; + } throw error; }); - return this.inviteCode; + const resumed = !! reply.resumed; + // A registration that did not reclaim the old session leaves every + // client of it unreachable, whatever the socket says. + for ( const channel of this.#channels.values() ) { + if ( resumed ) channel.onusable(); + else channel.strand(); + } + this.#resumeToken = reply.resumeToken ?? this.#resumeToken; + this.inviteCode = reply.invitecode ?? this.inviteCode; + return { inviteCode: this.inviteCode, resumed }; } - /** Settles the registration handshake, if one is still waiting on a reply. */ - #settleCreate ( settle ) { + #scheduleReconnect () { + if ( this.#closed || this.#reconnectTimer ) return; + const attempt = this.#reconnectAttempts++; + const backoff = Math.min(RECONNECT_MAX_MS, RECONNECT_BASE_MS * 2 ** attempt); + const delay = backoff / 2 + Math.random() * (backoff / 2); + this.#reconnectTimer = setTimeout(() => { + this.#reconnectTimer = null; + void this.#reconnect(); + }, delay); + } + + async #reconnect () { + if ( this.#closed ) return; + let registration; + try { + registration = await this.#register(); + } catch { + if ( ! this.#closed ) this.#scheduleReconnect(); + return; + } + this.#reconnectAttempts = 0; + this.dispatchEvent( + new PuterPeerServerReconnectEvent(registration.inviteCode, registration.resumed), + ); + } + + /** + * Settles the registration `ws` is waiting on, if it is still the one in + * flight - a replaced socket must not settle its successor's handshake. + * Passing no socket settles whichever is, which is what closing does. + */ + #settleCreate ( ws, settle ) { const pending = this.#pendingCreate; - if ( ! pending ) return; + if ( ! pending || ( ws && pending.ws !== ws ) ) return; + clearTimeout(pending.timer); this.#pendingCreate = null; settle(pending); } @@ -110,7 +199,7 @@ export class PuterPeerServer extends EventTarget { } } - async #message ( event ) { + async #message ( event, ws ) { let data; try { data = JSON.parse(event.data); @@ -121,9 +210,9 @@ export class PuterPeerServer extends EventTarget { if ( data.server.create ) { const reply = data.server.create; - this.#settleCreate((pending) => { + this.#settleCreate(ws, (pending) => { if ( reply.success ) { - pending.resolve(reply.invitecode); + pending.resolve(reply); } else { pending.reject(new Error(reply.error)); } @@ -161,9 +250,15 @@ export class PuterPeerServer extends EventTarget { } close () { - this.#settleCreate((pending) => { + this.#closed = true; + if ( this.#reconnectTimer ) clearTimeout(this.#reconnectTimer); + this.#reconnectTimer = null; + this.#settleCreate(null, (pending) => { pending.reject(new Error('The server was closed')); }); + // Give the invite code up rather than leave it held for a server + // that is never coming back. + if ( this.#alive ) this.relay({ release: {} }); for ( const connection of this.connections.values() ) connection.close(); this.connections.clear(); this.#channels.clear(); diff --git a/src/puter-js/src/modules/Peer/PuterPeerServer.test.js b/src/puter-js/src/modules/Peer/PuterPeerServer.test.js index a19fc1b5b..ec072f06e 100644 --- a/src/puter-js/src/modules/Peer/PuterPeerServer.test.js +++ b/src/puter-js/src/modules/Peer/PuterPeerServer.test.js @@ -28,14 +28,18 @@ class FakeWebSocket { } const origWebSocket = globalThis.WebSocket; +const origRTCPeerConnection = globalThis.RTCPeerConnection; beforeEach(() => { FakeWebSocket.latest = null; + FakePeerConnection.instances = []; globalThis.WebSocket = FakeWebSocket; + globalThis.RTCPeerConnection = FakePeerConnection; }); afterEach(() => { globalThis.WebSocket = origWebSocket; + globalThis.RTCPeerConnection = origRTCPeerConnection; vi.useRealTimers(); }); @@ -55,13 +59,31 @@ const startServer = async () => { await flushMicrotasks(); await FakeWebSocket.latest.onmessage({ data: JSON.stringify({ - server: { create: { success: true, invitecode: 'invite-1' } }, + server: { + create: { success: true, invitecode: 'invite-1', resumeToken: 'tok-1' }, + }, }), }); await started; return server; }; +/** The `create` request the most recent socket sent. */ +const createRequest = (ws) => + JSON.parse(ws.sent.find((raw) => JSON.parse(raw).server?.create)).server.create; + +/** Completes the open handshake, so the `create` request goes out. */ +const openSocket = async (ws) => { + ws.onopen(); + await flushMicrotasks(); +}; + +/** Answers the registration in flight on `ws`. */ +const answerCreate = async (ws, create) => { + await ws.onmessage({ data: JSON.stringify({ server: { create } }) }); + await flushMicrotasks(); +}; + describe('PuterPeerServer orphan offers', () => { it('ignores offers for unknown connection ids without throwing', async () => { await startServer(); @@ -183,3 +205,123 @@ describe('PuterPeerServer connections', () => { expect(ws.sent.length).toBe(before); }); }); + +describe('PuterPeerServer reclaiming a dropped session', () => { + it('presents its resume token and keeps the invite code it was reclaiming', async () => { + vi.useFakeTimers(); + const server = await startServer(); + const reconnects = []; + server.addEventListener('reconnect', (e) => reconnects.push([e.inviteCode, e.resumed])); + + FakeWebSocket.latest.onclose({}); + await vi.advanceTimersByTimeAsync(1_000); + + const second = FakeWebSocket.latest; + await openSocket(second); + expect(createRequest(second).resume).toBe('tok-1'); + + await answerCreate(second, { + success: true, + invitecode: 'invite-1', + resumeToken: 'tok-1', + resumed: true, + }); + + expect(server.inviteCode).toBe('invite-1'); + expect(server.signallingAlive).toBe(true); + expect(reconnects).toEqual([['invite-1', true]]); + }); + + it('takes the new code when the session could not be reclaimed', async () => { + vi.useFakeTimers(); + const server = await startServer(); + const reconnects = []; + server.addEventListener('reconnect', (e) => reconnects.push([e.inviteCode, e.resumed])); + + FakeWebSocket.latest.onclose({}); + await vi.advanceTimersByTimeAsync(1_000); + await openSocket(FakeWebSocket.latest); + await answerCreate(FakeWebSocket.latest, { + success: true, + invitecode: 'invite-2', + resumeToken: 'tok-2', + resumed: false, + }); + + expect(server.inviteCode).toBe('invite-2'); + expect(reconnects).toEqual([['invite-2', false]]); + }); + + it('strands the connections a fresh registration left unroutable', async () => { + // The signaller has no route to them any more, so a link that kept + // trying would only wait out every timeout it has before dying. + vi.useFakeTimers(); + const server = await startServer(); + const ws = FakeWebSocket.latest; + await ws.onmessage({ + data: JSON.stringify({ server: { connect: { id: 'c1', user: {} } } }), + }); + const conn = server.connections.get('c1'); + const pc = FakePeerConnection.instances.at(-1); + // Up and carrying traffic: a link that never came up closes on the + // socket alone, which is not what this is about. + pc.channels[0].open(); + const closes = []; + conn.addEventListener('close', (e) => closes.push(e.reason)); + + ws.onclose({}); + await vi.advanceTimersByTimeAsync(1_000); + await openSocket(FakeWebSocket.latest); + await answerCreate(FakeWebSocket.latest, { + success: true, + invitecode: 'invite-2', + resumeToken: 'tok-2', + resumed: false, + }); + + pc.setConnectionState('failed'); + await flushMicrotasks(); + + expect(pc.restarts).toBe(0); + expect(conn.closed).toBe(true); + expect(closes).toEqual(['the peer is no longer reachable']); + }); + + it('keeps trying while the signaller is unreachable', async () => { + vi.useFakeTimers(); + const server = await startServer(); + + FakeWebSocket.latest.onclose({}); + await vi.advanceTimersByTimeAsync(1_000); + const second = FakeWebSocket.latest; + await openSocket(second); + second.onclose({}); + await vi.advanceTimersByTimeAsync(4_000); + + expect(FakeWebSocket.latest).not.toBe(second); + expect(server.signallingAlive).toBe(false); + }); + + it('gives the invite code up when the server is closed for good', async () => { + const server = await startServer(); + const ws = FakeWebSocket.latest; + const before = ws.sent.length; + + server.close(); + + const released = ws.sent.slice(before).map((raw) => JSON.parse(raw)); + expect(released).toContainEqual({ server: { release: {} } }); + }); + + it('does not reconnect after close()', async () => { + vi.useFakeTimers(); + const server = await startServer(); + const ws = FakeWebSocket.latest; + + server.close(); + ws.onclose?.({}); + await vi.advanceTimersByTimeAsync(10_000); + + expect(FakeWebSocket.latest).toBe(ws); + }); +}); diff --git a/src/puter-js/src/modules/Peer/events.js b/src/puter-js/src/modules/Peer/events.js index e6aa2e24c..8f72baf38 100644 --- a/src/puter-js/src/modules/Peer/events.js +++ b/src/puter-js/src/modules/Peer/events.js @@ -8,6 +8,21 @@ export class PuterPeerServerConnectionEvent extends Event { } } +/** + * The signaller socket came back. `resumed` says whether the session came + * with it: when it did, the invite code is the same one and the clients + * already connected can be renegotiated with again. + */ +export class PuterPeerServerReconnectEvent extends Event { + inviteCode; + resumed; + constructor (inviteCode, resumed) { + super('reconnect'); + this.inviteCode = inviteCode; + this.resumed = resumed; + } +} + export class PuterPeerConnectionMessageEvent extends Event { data; constructor (message) { diff --git a/src/puter-js/src/modules/Peer/signalling.js b/src/puter-js/src/modules/Peer/signalling.js index 002cad97e..34202ebb1 100644 --- a/src/puter-js/src/modules/Peer/signalling.js +++ b/src/puter-js/src/modules/Peer/signalling.js @@ -20,12 +20,19 @@ export class SignallingChannel { oncandidate = () => {}; /** @type {(reason?: string) => void} */ onbye = () => {}; - /** Peer's signalling session ended. Evidence it hung up, never proof. */ - /** @type {(reason?: string) => void} */ + /** + * Peer's signalling session ended. Evidence it hung up, never proof - + * and when the signaller says the session is reclaimable, not even + * that: the peer is expected back on it. + */ + /** @type {(reason?: string, resumable?: boolean) => void} */ onpeergone = () => {}; /** Our own signalling path died; renegotiation is unavailable. */ /** @type {() => void} */ onunusable = () => {}; + /** It came back, and anything waiting on it may go ahead. */ + /** @type {() => void} */ + onusable = () => {}; /** * Reads one relayed envelope and calls the handler for what it holds. @@ -164,7 +171,7 @@ export class ClientSignallingChannel extends SignallingChannel { return this.onrejected(new Error(msg.connect.error)); } if ( msg.disconnect ) { - this.onpeergone(msg.disconnect.reason); + this.onpeergone(msg.disconnect.reason, !! msg.disconnect.resumable); return; } this.receive(msg); @@ -184,6 +191,13 @@ export class ClientSignallingChannel extends SignallingChannel { export class ServerSignallingChannel extends SignallingChannel { #server; #id; + /** + * Set when the server re-registered without reclaiming its session. The + * signaller has no route to this client any more and never will, so the + * connection is better off hearing that at once than waiting out every + * timeout it has. + */ + #stranded = false; /** * @param {import('./PuterPeerServer.js').PuterPeerServer} server @@ -196,7 +210,18 @@ export class ServerSignallingChannel extends SignallingChannel { } get alive () { - return this.#server.signallingAlive; + return ! this.#stranded && this.#server.signallingAlive; + } + + get stranded () { + return this.#stranded; + } + + /** @returns {void} */ + strand () { + if ( this.#stranded ) return; + this.#stranded = true; + this.onunusable(); } // Every payload a peer server sends carries the connection id, since one