mirror of
https://github.com/HeyPuter/puter.git
synced 2026-10-11 06:11:56 +00:00
proper reconnect
This commit is contained in:
1 parent
0ced965207
commit
75e437013f
6 files changed
+448
-74
No files matched your search
@@ -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<void>}
|
||||
*/
|
||||
#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 ) {
|
||||
|
||||
@@ -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');
|
||||
});
|
||||
|
||||
@@ -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<string>}
|
||||
*/
|
||||
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();
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in new issue
Block a user