Files
Toju/toju-app/src/app/infrastructure/realtime/peer-connection-manager/recovery/peer-recovery.ts
T
myxelium a83f5aa750 fix(realtime): survive signal outages and per-server identities
Peer recovery burned its whole retry budget while signaling was down, then
dropped the tracker with no re-arm, so a peer stayed dead until an unrelated
roster event healed it. Recovery now waits for a usable transport before
spending an attempt.

Initiator election also compared a home actor id against foreign roster ids,
which is not antisymmetric across signal servers - both sides offered, or
neither did. Election moves into `peer-role.rules` and compares ids only
within one signal server's identity space.
2026-08-14 03:19:29 +02:00

563 lines
16 KiB
TypeScript

import {
CONNECTION_STATE_CONNECTED,
DATA_CHANNEL_RECOVERY_GRACE_MS,
DATA_CHANNEL_STATE_OPEN,
P2P_TYPE_VOICE_STATE_REQUEST,
PEER_DISCONNECT_GRACE_MS,
PEER_RECONNECT_INTERVAL_MS,
PEER_RECONNECT_MAX_ATTEMPTS
} from '../../realtime.constants';
import { shouldInitiatePeerConnection } from '../../peer-role.rules';
import {
PeerConnectionManagerContext,
PeerConnectionManagerState,
RecoveryHandlers,
RemovePeerOptions
} from '../shared';
import { clearAllPingTimers, stopPingInterval } from '../messaging/ping';
import { clearDebugNetworkPeerMetrics } from '../../logging/debug-network-metrics';
/**
* Close and remove a peer connection, data channel, and emit a disconnect event.
*/
export function removePeer(
context: PeerConnectionManagerContext,
peerId: string,
options?: RemovePeerOptions
): void {
const { state } = context;
const peerData = state.activePeerConnections.get(peerId);
const preserveReconnectState = options?.preserveReconnectState === true;
clearPeerDisconnectGraceTimer(state, peerId);
clearDataChannelRecoveryTimer(state, peerId);
if (!preserveReconnectState) {
clearPeerReconnectTimer(state, peerId);
state.disconnectedPeerTracker.delete(peerId);
clearDebugNetworkPeerMetrics(peerId);
}
state.remotePeerStreams.delete(peerId);
state.remotePeerVoiceStreams.delete(peerId);
state.remotePeerScreenShareStreams.delete(peerId);
state.remotePeerCameraStreams.delete(peerId);
if (peerData) {
if (peerData.dataChannel)
peerData.dataChannel.close();
peerData.connection.close();
state.activePeerConnections.delete(peerId);
state.peerNegotiationQueue.delete(peerId);
removeFromConnectedPeers(state, peerId);
stopPingInterval(state, peerId);
state.peerLatencies.delete(peerId);
state.pendingPings.delete(peerId);
state.peerDisconnected$.next(peerId);
}
}
/** Close every active peer connection and clear internal state. */
export function closeAllPeers(state: PeerConnectionManagerState): void {
clearAllPeerReconnectTimers(state);
clearAllPeerDisconnectGraceTimers(state);
clearAllDataChannelRecoveryTimers(state);
clearAllPingTimers(state);
state.activePeerConnections.forEach((peerData) => {
if (peerData.dataChannel)
peerData.dataChannel.close();
peerData.connection.close();
});
state.activePeerConnections.clear();
state.remotePeerStreams.clear();
state.remotePeerVoiceStreams.clear();
state.remotePeerScreenShareStreams.clear();
state.remotePeerCameraStreams.clear();
state.peerNegotiationQueue.clear();
state.peerLatencies.clear();
state.pendingPings.clear();
state.connectedPeersChanged$.next([]);
}
export function trackDisconnectedPeer(state: PeerConnectionManagerState, peerId: string): void {
state.disconnectedPeerTracker.set(peerId, {
lastSeenTimestamp: Date.now(),
reconnectAttempts: 0
});
}
export function clearPeerReconnectTimer(
state: PeerConnectionManagerState,
peerId: string
): void {
const timer = state.peerReconnectTimers.get(peerId);
if (timer) {
clearInterval(timer);
state.peerReconnectTimers.delete(peerId);
}
}
export function clearPeerDisconnectGraceTimer(
state: PeerConnectionManagerState,
peerId: string
): void {
const timer = state.peerDisconnectGraceTimers.get(peerId);
if (timer) {
clearTimeout(timer);
state.peerDisconnectGraceTimers.delete(peerId);
}
}
export function clearDataChannelRecoveryTimer(
state: PeerConnectionManagerState,
peerId: string
): void {
const timer = state.dataChannelRecoveryTimers.get(peerId);
if (timer) {
clearTimeout(timer);
state.dataChannelRecoveryTimers.delete(peerId);
}
}
/** Cancel all pending peer reconnect timers and clear the tracker. */
export function clearAllPeerReconnectTimers(state: PeerConnectionManagerState): void {
state.peerReconnectTimers.forEach((timer) => clearInterval(timer));
state.peerReconnectTimers.clear();
state.disconnectedPeerTracker.clear();
}
export function clearAllPeerDisconnectGraceTimers(state: PeerConnectionManagerState): void {
state.peerDisconnectGraceTimers.forEach((timer) => clearTimeout(timer));
state.peerDisconnectGraceTimers.clear();
}
export function clearAllDataChannelRecoveryTimers(state: PeerConnectionManagerState): void {
state.dataChannelRecoveryTimers.forEach((timer) => clearTimeout(timer));
state.dataChannelRecoveryTimers.clear();
}
export function scheduleDataChannelRecovery(
context: PeerConnectionManagerContext,
peerId: string,
channel: RTCDataChannel,
reason: string,
handlers: RecoveryHandlers
): void {
const { logger, state } = context;
const peerData = state.activePeerConnections.get(peerId);
if (!peerData || peerData.dataChannel !== channel)
return;
if (channel.readyState === DATA_CHANNEL_STATE_OPEN)
return;
if (channel.readyState === 'closed') {
logger.warn('[data-channel] Control channel closed; reconnecting peer immediately', {
channelLabel: channel.label,
connectionState: peerData.connection.connectionState,
peerId,
reason
});
repairUnavailableDataChannel(context, peerId, channel, reason, handlers);
return;
}
if (state.dataChannelRecoveryTimers.has(peerId))
return;
logger.warn('[data-channel] Control channel unavailable; waiting before reconnect', {
channelLabel: channel.label,
peerId,
readyState: channel.readyState,
reason
});
const timer = setTimeout(() => {
state.dataChannelRecoveryTimers.delete(peerId);
const latestPeerData = state.activePeerConnections.get(peerId);
if (!latestPeerData || latestPeerData.dataChannel !== channel)
return;
if (latestPeerData.dataChannel?.readyState === DATA_CHANNEL_STATE_OPEN)
return;
logger.warn('[data-channel] Control channel did not recover; selecting repair path', {
channelLabel: channel.label,
connectionState: latestPeerData.connection.connectionState,
peerId,
readyState: latestPeerData.dataChannel?.readyState ?? null,
reason
});
repairUnavailableDataChannel(context, peerId, channel, reason, handlers);
}, DATA_CHANNEL_RECOVERY_GRACE_MS);
state.dataChannelRecoveryTimers.set(peerId, timer);
}
/**
* Repair a failed control channel with the least destructive option available.
*
* While the `RTCPeerConnection` is still connected the channel is replaced on that same
* connection, so voice, camera, and screen share keep flowing. Only the deterministically
* elected initiator creates the replacement; the other side adopts the incoming channel.
* A full peer rebuild is the fallback when the connection itself is gone, when the
* replacement cannot be created, or when the replacement never opens.
*/
function repairUnavailableDataChannel(
context: PeerConnectionManagerContext,
peerId: string,
channel: RTCDataChannel,
reason: string,
handlers: RecoveryHandlers
): void {
const { callbacks, logger, state } = context;
const peerData = state.activePeerConnections.get(peerId);
if (!peerData || peerData.dataChannel !== channel)
return;
if (peerData.dataChannel?.readyState === DATA_CHANNEL_STATE_OPEN)
return;
if (peerData.connection.connectionState !== CONNECTION_STATE_CONNECTED) {
rebuildPeerAfterDataChannelFailure(context, peerId, reason, handlers);
return;
}
const localOderId = callbacks.getIdentifyCredentialsForPeer(peerId)?.oderId ?? null;
if (!localOderId) {
logger.warn('[data-channel] Logical identity unknown; rebuilding peer instead of replacing channel', {
peerId,
reason
});
rebuildPeerAfterDataChannelFailure(context, peerId, reason, handlers);
return;
}
if (shouldInitiatePeerConnection(localOderId, peerId)) {
logger.warn('[data-channel] Replacing control channel on the live connection', {
channelLabel: channel.label,
peerId,
reason
});
if (!handlers.replaceDataChannel(peerId, channel)) {
rebuildPeerAfterDataChannelFailure(context, peerId, reason, handlers);
return;
}
watchDataChannelReplacement(context, peerId, reason, handlers);
return;
}
logger.info('[data-channel] Waiting for the remote replacement control channel', {
channelLabel: channel.label,
localOderId,
peerId,
reason
});
watchDataChannelReplacement(context, peerId, reason, handlers);
}
/** Rebuild the peer if the replacement control channel does not open within the grace period. */
function watchDataChannelReplacement(
context: PeerConnectionManagerContext,
peerId: string,
reason: string,
handlers: RecoveryHandlers
): void {
const { logger, state } = context;
clearDataChannelRecoveryTimer(state, peerId);
const timer = setTimeout(() => {
state.dataChannelRecoveryTimers.delete(peerId);
const latestPeerData = state.activePeerConnections.get(peerId);
if (!latestPeerData)
return;
if (latestPeerData.dataChannel?.readyState === DATA_CHANNEL_STATE_OPEN) {
logger.info('[data-channel] Replacement control channel is open; peer kept alive', {
peerId,
reason
});
return;
}
logger.warn('[data-channel] Replacement control channel never opened; rebuilding peer', {
connectionState: latestPeerData.connection.connectionState,
peerId,
readyState: latestPeerData.dataChannel?.readyState ?? null,
reason
});
rebuildPeerAfterDataChannelFailure(context, peerId, reason, handlers);
}, DATA_CHANNEL_RECOVERY_GRACE_MS);
state.dataChannelRecoveryTimers.set(peerId, timer);
}
function rebuildPeerAfterDataChannelFailure(
context: PeerConnectionManagerContext,
peerId: string,
reason: string,
handlers: RecoveryHandlers
): void {
const { logger, state } = context;
const peerData = state.activePeerConnections.get(peerId);
logger.warn('[data-channel] Recreating peer transport after control channel failure', {
connectionState: peerData?.connection.connectionState ?? null,
peerId,
readyState: peerData?.dataChannel?.readyState ?? null,
reason
});
trackDisconnectedPeer(state, peerId);
handlers.removePeer(peerId, { preserveReconnectState: true });
attemptPeerReconnect(context, peerId, handlers);
schedulePeerReconnect(context, peerId, handlers);
}
export function schedulePeerDisconnectRecovery(
context: PeerConnectionManagerContext,
peerId: string,
handlers: RecoveryHandlers
): void {
const { logger, state } = context;
if (state.peerDisconnectGraceTimers.has(peerId))
return;
logger.warn('Peer temporarily disconnected; waiting before reconnect', { peerId });
const timer = setTimeout(() => {
state.peerDisconnectGraceTimers.delete(peerId);
const peerData = state.activePeerConnections.get(peerId);
if (!peerData)
return;
const connectionState = peerData.connection.connectionState;
if (connectionState === CONNECTION_STATE_CONNECTED || connectionState === 'connecting') {
logger.info('Peer recovered before disconnect grace expired', {
peerId,
state: connectionState
});
return;
}
logger.warn('Peer still disconnected after grace period; recreating connection', {
peerId,
state: connectionState
});
trackDisconnectedPeer(state, peerId);
handlers.removePeer(peerId, { preserveReconnectState: true });
schedulePeerReconnect(context, peerId, handlers);
}, PEER_DISCONNECT_GRACE_MS);
state.peerDisconnectGraceTimers.set(peerId, timer);
}
export function schedulePeerReconnect(
context: PeerConnectionManagerContext,
peerId: string,
handlers: RecoveryHandlers
): void {
const { callbacks, logger, state } = context;
if (state.peerReconnectTimers.has(peerId))
return;
logger.info('Scheduling P2P reconnect', { peerId });
const timer = setInterval(() => {
const info = state.disconnectedPeerTracker.get(peerId);
if (!info) {
clearPeerReconnectTimer(state, peerId);
return;
}
// An attempt costs nothing while the socket is down and there is no way to send an
// offer, so the budget must only be spent on attempts that can actually reach the peer.
if (!callbacks.isSignalingConnected()) {
logger.info('Deferring P2P reconnect - no signaling connection', {
peerId,
attemptsSpent: info.reconnectAttempts
});
return;
}
if (info.reconnectAttempts >= PEER_RECONNECT_MAX_ATTEMPTS) {
failPeerRecovery(context, peerId, info.reconnectAttempts);
return;
}
info.reconnectAttempts++;
logger.info('P2P reconnect attempt', {
peerId,
attempt: info.reconnectAttempts
});
attemptPeerReconnect(context, peerId, handlers);
}, PEER_RECONNECT_INTERVAL_MS);
state.peerReconnectTimers.set(peerId, timer);
}
/**
* Stop retrying a peer and say so. The tracker entry is kept (marked failed) so
* `resumeStalledPeerRecovery` can re-arm it once signaling is usable again.
*/
function failPeerRecovery(
context: PeerConnectionManagerContext,
peerId: string,
attempts: number
): void {
const { logger, state } = context;
const info = state.disconnectedPeerTracker.get(peerId);
logger.warn('P2P reconnect gave up', { attempts, peerId });
clearPeerReconnectTimer(state, peerId);
if (info) {
info.recoveryFailed = true;
}
state.peerRecoveryStatus$.next({ attempts, peerId, status: 'failed' });
}
/**
* Re-arm every peer whose recovery gave up. Called when signaling reconnects, so a peer
* that ran out of attempts during an outage is retried instead of staying dead silently.
*/
export function resumeStalledPeerRecovery(
context: PeerConnectionManagerContext,
handlers: RecoveryHandlers
): void {
const { callbacks, logger, state } = context;
if (!callbacks.isSignalingConnected())
return;
state.disconnectedPeerTracker.forEach((info, peerId) => {
if (!info.recoveryFailed)
return;
logger.info('Re-arming P2P recovery after signaling returned', { peerId });
info.recoveryFailed = false;
info.reconnectAttempts = 0;
state.peerRecoveryStatus$.next({ attempts: 0, peerId, status: 'retrying' });
attemptPeerReconnect(context, peerId, handlers);
schedulePeerReconnect(context, peerId, handlers);
});
}
export function attemptPeerReconnect(
context: PeerConnectionManagerContext,
peerId: string,
handlers: RecoveryHandlers
): void {
const { callbacks, logger, state } = context;
if (state.activePeerConnections.has(peerId)) {
handlers.removePeer(peerId, { preserveReconnectState: true });
}
const localOderId = callbacks.getIdentifyCredentialsForPeer(peerId)?.oderId ?? null;
if (!localOderId) {
logger.info('Skipping reconnect offer until logical identity is ready', { peerId });
handlers.createPeerConnection(peerId, false);
return;
}
const shouldInitiate = shouldInitiatePeerConnection(localOderId, peerId);
handlers.createPeerConnection(peerId, shouldInitiate);
if (shouldInitiate) {
void handlers.createAndSendOffer(peerId);
return;
}
logger.info('Waiting for remote reconnect offer based on deterministic initiator selection', {
localOderId,
peerId
});
}
export function requestVoiceStateFromPeer(
state: PeerConnectionManagerState,
logger: PeerConnectionManagerContext['logger'],
peerId: string
): void {
const peerData = state.activePeerConnections.get(peerId);
if (peerData?.dataChannel?.readyState === DATA_CHANNEL_STATE_OPEN) {
try {
peerData.dataChannel.send(JSON.stringify({ type: P2P_TYPE_VOICE_STATE_REQUEST }));
} catch (error) {
logger.warn('Failed to request voice state', error);
}
}
}
/** Return a snapshot copy of the currently-connected peer IDs. */
export function getConnectedPeerIds(state: PeerConnectionManagerState): string[] {
return [...state.connectedPeersList];
}
export function addToConnectedPeers(state: PeerConnectionManagerState, peerId: string): void {
if (!state.connectedPeersList.includes(peerId)) {
state.connectedPeersList = [...state.connectedPeersList, peerId];
state.connectedPeersChanged$.next(state.connectedPeersList);
}
}
/**
* Remove a peer from the connected list and notify subscribers.
*/
export function removeFromConnectedPeers(
state: PeerConnectionManagerState,
peerId: string
): void {
state.connectedPeersList = state.connectedPeersList.filter(
(connectedId) => connectedId !== peerId
);
state.connectedPeersChanged$.next(state.connectedPeersList);
}
/** Reset the connected peers list to empty and notify subscribers. */
export function resetConnectedPeers(state: PeerConnectionManagerState): void {
state.connectedPeersList = [];
state.connectedPeersChanged$.next([]);
}