import type { PacketChannel } from './channels'; import type { Envelope } from './protocol'; interface DataChannelLike { readonly readyState: string; readonly bufferedAmount?: number; send(data: string): void; close(): void; addEventListener(type: 'message', handler: (event: { data: unknown }) => void): void; addEventListener(type: 'close', handler: () => void): void; } type EnvelopeHandler = (envelope: Envelope) => void; type CloseHandler = (reason?: string) => void; /** * The label of the SECOND data channel — unordered, `maxRetransmits: 0` — that * carries `sendUnreliable` traffic. The offerer creates it alongside the * reliable channel; the answerer receives it via `ondatachannel` and hands it * to `attachUnreliableChannel`. Exported so the negotiation code (`engine.ts`) * and this class agree on the one string. */ export const WEBRTC_UNRELIABLE_CHANNEL_LABEL = 'p2p-colyseus-unreliable'; export class WebRTCDataChannelPacketChannel implements PacketChannel { readonly mode = 'webrtc' as const; private readonly envelopeHandlers = new Set(); private readonly closeHandlers = new Set(); private unreliableChannel: DataChannelLike | null = null; constructor( readonly peerId: string, private readonly dataChannel: DataChannelLike, unreliableChannel?: DataChannelLike, ) { this.bindReceive(dataChannel); dataChannel.addEventListener('close', () => { for (const handler of this.closeHandlers) handler('closed'); }); if (unreliableChannel) this.attachUnreliableChannel(unreliableChannel); } /** * Attach the unreliable sub-channel after construction. The answerer learns * of it through a SEPARATE `ondatachannel` event, which may arrive before or * after the reliable channel that constructs this object; the offerer passes * it in the constructor instead. Frames received on it dispatch through the * same envelope handlers as the reliable channel. The unreliable channel does * NOT govern this channel's lifecycle — the reliable channel's `close` does — * so no close handler is bound to it. */ attachUnreliableChannel(channel: DataChannelLike): void { this.unreliableChannel = channel; this.bindReceive(channel); } private bindReceive(channel: DataChannelLike): void { channel.addEventListener('message', (event) => { if (typeof event.data !== 'string') return; const envelope = JSON.parse(event.data) as Envelope; for (const handler of this.envelopeHandlers) handler(envelope); }); } send(envelope: Envelope): void { if (this.dataChannel.readyState !== 'open') { throw new Error(`WebRTC data channel for ${this.peerId} is not open`); } this.dataChannel.send(JSON.stringify(envelope)); } sendUnreliable(envelope: Envelope): void { const channel = this.unreliableChannel; if (channel && channel.readyState === 'open') { channel.send(JSON.stringify(envelope)); return; } // No unreliable channel yet (still negotiating) or it has gone away — never // drop or throw: fall back to the reliable channel so the message arrives. this.send(envelope); } onEnvelope(handler: EnvelopeHandler): () => void { this.envelopeHandlers.add(handler); return () => this.envelopeHandlers.delete(handler); } onClose(handler: CloseHandler): () => void { this.closeHandlers.add(handler); return () => this.closeHandlers.delete(handler); } close(): void { this.dataChannel.close(); this.unreliableChannel?.close(); } closeAfterFlush(delayMs = 100): void { const started = Date.now(); const minDelayMs = delayMs; const maxDelayMs = Math.max(1_000, delayMs); const tick = () => { const elapsed = Date.now() - started; const flushed = (this.dataChannel.bufferedAmount ?? 0) === 0; if ((flushed && elapsed >= minDelayMs) || elapsed >= maxDelayMs) { this.close(); return; } setTimeout(tick, 10); }; tick(); } } export function createWebRTCDataChannelPacketChannel( peerId: string, dataChannel: RTCDataChannel, unreliableChannel?: RTCDataChannel, ): PacketChannel { return new WebRTCDataChannelPacketChannel(peerId, dataChannel, unreliableChannel); }