import type { RelayProofDiagnostics } from '../diagnostics'; import { P2P_CLOSE_CODES } from '../protocol'; import { DEFAULT_RELAY_LIMITS, type RelayLimits } from './limits'; import type { SignalEnvelope } from './protocol'; export interface RelayLogger { info(event: string, fields?: Record): void; } interface SignalPeer { id: string; sent: SignalEnvelope[]; handlers: Set<(envelope: SignalEnvelope) => void>; } interface RateWindow { startedAt: number; count: number; } interface PrincipalQuotaRecord { hostedRoomIds: Set; relayRoomIds: Set; relayBytesByDay: Map; } /** * Quota store keyed by a neutral `principal` (an account id, or — as dormant * scaffolding — an ip hash). The relay's durable per-account byte allotment * lives in Durable Object storage (see `durable-room.ts`); this in-memory store * backs the coordinator's per-principal concurrent-room and per-day byte caps. */ export interface RelayQuotaStore { reserveHostedRoom(principal: string, roomId: string, maxRooms: number): boolean; reserveRelayRoom(principal: string, roomId: string, maxRooms: number): boolean; consumeRelayBytes( principal: string, dayKey: string, bytes: number, maxBytesPerDay: number, ): boolean; } export class InMemoryRelayQuotaStore implements RelayQuotaStore { private readonly records = new Map(); reserveHostedRoom(principal: string, roomId: string, maxRooms: number): boolean { const record = this.recordFor(principal); record.hostedRoomIds.add(roomId); return record.hostedRoomIds.size <= maxRooms; } reserveRelayRoom(principal: string, roomId: string, maxRooms: number): boolean { const record = this.recordFor(principal); record.relayRoomIds.add(roomId); return record.relayRoomIds.size <= maxRooms; } consumeRelayBytes( principal: string, dayKey: string, bytes: number, maxBytesPerDay: number, ): boolean { const record = this.recordFor(principal); const next = (record.relayBytesByDay.get(dayKey) ?? 0) + bytes; if (next > maxBytesPerDay) return false; record.relayBytesByDay.set(dayKey, next); return true; } private recordFor(principal: string): PrincipalQuotaRecord { let record = this.records.get(principal); if (!record) { record = { hostedRoomIds: new Set(), relayRoomIds: new Set(), relayBytesByDay: new Map() }; this.records.set(principal, record); } return record; } } export class SignalingRelayCoordinator { private roomId = ''; private roomName = ''; private hostSignalId = ''; private createdAt = 0; private lastActivityAt = 0; private lastHostHeartbeatAt = 0; private readonly peers = new Map(); private readonly signalRate = new Map(); private readonly relayRate = new Map(); private readonly iceCandidates = new Map(); private readonly relayAllowed = new Set(); private offerCount = 0; private answerCount = 0; private iceCandidateCount = 0; private relayOpenCount = 0; private relayBytesIn = 0; private relayBytesOut = 0; constructor( private readonly limits: RelayLimits = DEFAULT_RELAY_LIMITS, private readonly now: () => number = () => Date.now(), private readonly quotaStore: RelayQuotaStore = new InMemoryRelayQuotaStore(), private readonly logger?: RelayLogger | undefined, ) {} /** * The room name bound at host-register (empty until a host registers). The * POST relay-data biller verifies each managed access token against it, so a * token minted for a different room can never meter bytes here. */ get boundRoomName(): string { return this.roomName; } registerHost(roomName: string, hostSignalId = 'host', principal?: string): SignalEnvelope { const rateLimit = this.consumeSignal(hostSignalId); if (rateLimit) return rateLimit; const now = this.now(); this.roomId = crypto.randomUUID?.() ?? `room-${Date.now()}`; if ( principal && !this.quotaStore.reserveHostedRoom( principal, this.roomId, this.limits.maxConcurrentHostedRoomsPerIpHash, ) ) { this.roomId = ''; return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } this.roomName = roomName; this.hostSignalId = hostSignalId; this.createdAt = now; this.lastActivityAt = now; this.lastHostHeartbeatAt = now; this.peers.clear(); this.signalRate.clear(); this.relayRate.clear(); this.iceCandidates.clear(); this.relayAllowed.clear(); this.offerCount = 0; this.answerCount = 0; this.iceCandidateCount = 0; this.relayOpenCount = 0; this.relayBytesIn = 0; this.relayBytesOut = 0; this.peers.set(hostSignalId, { id: hostSignalId, sent: [], handlers: new Set() }); this.log('relay.host_registered', { roomId: this.roomId, roomName, hostSignalId }); return { kind: 'host-registered', roomId: this.roomId, hostSignalId }; } ensureRelayPeer(roomId: string, peerId: string, targetPeerId: string): void { const now = this.now(); if (!this.roomId) { this.roomId = roomId; this.roomName = 'websocket-relay'; this.hostSignalId = peerId; this.createdAt = now; this.lastActivityAt = now; this.lastHostHeartbeatAt = now; } if (this.roomId !== roomId) return; this.peers.set(peerId, this.peers.get(peerId) ?? { id: peerId, sent: [], handlers: new Set() }); this.peers.set( targetPeerId, this.peers.get(targetPeerId) ?? { id: targetPeerId, sent: [], handlers: new Set() }, ); } join(roomName: string, clientSignalId: string): SignalEnvelope { const rateLimit = this.consumeSignal(clientSignalId); if (rateLimit) return rateLimit; const unavailable = this.unavailableRoomError(); if (unavailable) return unavailable; if (!this.roomId || roomName !== this.roomName) { return { kind: 'error', message: String(P2P_CLOSE_CODES.roomNotFound) }; } if (this.peers.size >= this.limits.maxPlayers) { return { kind: 'error', message: String(P2P_CLOSE_CODES.roomFull) }; } this.lastActivityAt = this.now(); this.peers.set(clientSignalId, { id: clientSignalId, sent: [], handlers: new Set() }); this.log('relay.join_routed', { roomId: this.roomId, roomName, clientSignalId }); return { kind: 'join-routed', roomId: this.roomId, hostSignalId: this.hostSignalId, clientSignalId, }; } forward(envelope: SignalEnvelope & { target?: string }): SignalEnvelope | null { const sourcePeerId = sourceOf(envelope); if (envelope.kind !== 'relay-data') { const signalRateLimit = this.consumeSignal(sourcePeerId); if (signalRateLimit) return signalRateLimit; } const unavailable = this.unavailableRoomError(); if (unavailable) return unavailable; if (envelope.kind === 'rtc-offer') this.offerCount++; if (envelope.kind === 'rtc-answer') this.answerCount++; if (envelope.kind === 'rtc-ice') { const count = (this.iceCandidates.get(sourcePeerId) ?? 0) + 1; this.iceCandidates.set(sourcePeerId, count); if (count > this.limits.maxIceCandidatesPerPeerPerJoin) { return { kind: 'error', message: String(P2P_CLOSE_CODES.rateLimited) }; } this.iceCandidateCount++; } if (envelope.kind === 'relay-open') { if ( envelope.ipHash && !this.quotaStore.reserveRelayRoom( envelope.ipHash, envelope.roomId, this.limits.maxConcurrentRelayRoomsPerIpHash, ) ) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } this.relayAllowed.add(sourcePeerId); this.relayOpenCount++; this.log('relay.opened', { roomId: envelope.roomId, from: sourcePeerId, target: envelope.target, }); } if (envelope.kind === 'relay-data') { if (!this.relayAllowed.has(sourcePeerId)) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayNotAllowed) }; } const relayRateLimit = this.consumeRelay(sourcePeerId); if (relayRateLimit) return relayRateLimit; const bytes = byteLength(envelope); if (bytes > this.limits.maxEnvelopeBytes) { return { kind: 'error', message: String(P2P_CLOSE_CODES.envelopeTooLarge) }; } if (this.relayBytesIn + bytes > this.limits.maxRelayBytesPerRoom) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } if ( envelope.ipHash && !this.quotaStore.consumeRelayBytes( envelope.ipHash, dayKey(this.now()), bytes, this.limits.maxRelayBytesPerIpHashPerDay, ) ) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } this.relayBytesIn += bytes; this.relayBytesOut += bytes; this.log('relay.data', { roomId: envelope.roomId, from: sourcePeerId, target: envelope.target, bytes, }); } const target = 'target' in envelope ? envelope.target : undefined; if (!target) return null; const peer = this.peers.get(target); this.lastActivityAt = this.now(); peer?.sent.push(envelope); for (const handler of peer?.handlers ?? []) handler(envelope); return null; } onEnvelope(peerId: string, handler: (envelope: SignalEnvelope) => void): () => void { const peer = this.peers.get(peerId); if (!peer) throw new Error(`Unknown signaling peer: ${peerId}`); peer.handlers.add(handler); return () => peer.handlers.delete(handler); } drain(peerId: string): SignalEnvelope[] { const peer = this.peers.get(peerId); if (!peer) return []; this.lastActivityAt = this.now(); const sent = peer.sent.slice(); peer.sent.length = 0; return sent; } heartbeat(roomId: string, peerId: string): SignalEnvelope { const rateLimit = this.consumeSignal(peerId); if (rateLimit) return rateLimit; if (roomId !== this.roomId) return { kind: 'error', message: String(P2P_CLOSE_CODES.roomNotFound) }; if (peerId === this.hostSignalId) this.lastHostHeartbeatAt = this.now(); this.lastActivityAt = this.now(); this.log('relay.heartbeat', { roomId, peerId }); return { kind: 'heartbeat', roomId }; } diagnostics(): RelayProofDiagnostics { const maybeBudgetLimits = this.limits as RelayLimits & { maxGlobalRequestsPerDay?: number; maxGlobalWebSocketMessagesPerDay?: number; maxActiveRelaySockets?: number; }; return { roomId: this.roomId, hostSignalId: this.hostSignalId, clientSignalIds: [...this.peers.keys()].filter((id) => id !== this.hostSignalId), offerCount: this.offerCount, answerCount: this.answerCount, iceCandidateCount: this.iceCandidateCount, relayOpenCount: this.relayOpenCount, relayBytesIn: this.relayBytesIn, relayBytesOut: this.relayBytesOut, authoritativeStateOwned: false, quota: { maxPlayers: this.limits.maxPlayers, maxEnvelopeBytes: this.limits.maxEnvelopeBytes, maxMessagesPerSecond: this.limits.maxMessagesPerPeerPerSecond, maxGlobalRelayBytesPerDay: this.limits.maxGlobalRelayBytesPerDay, maxGlobalRequestsPerDay: maybeBudgetLimits.maxGlobalRequestsPerDay ?? 250_000, maxGlobalWebSocketMessagesPerDay: maybeBudgetLimits.maxGlobalWebSocketMessagesPerDay ?? 2_000_000, maxActiveRelaySockets: maybeBudgetLimits.maxActiveRelaySockets ?? 200, maxRelayBytesPerRoom: this.limits.maxRelayBytesPerRoom, }, }; } private consumeSignal(peerId: string): SignalEnvelope | null { return consumeRate( this.signalRate, peerId, this.now(), 60_000, this.limits.maxSignalingMessagesPerPeerPerMinute, ); } private consumeRelay(peerId: string): SignalEnvelope | null { return consumeRate( this.relayRate, peerId, this.now(), 1_000, this.limits.maxMessagesPerPeerPerSecond, ); } private unavailableRoomError(): SignalEnvelope | null { if (!this.roomId) return null; const now = this.now(); if (now - this.createdAt > this.limits.maxRelayRoomDurationMs) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } if (now - this.lastActivityAt > this.limits.idleTimeoutMs) { return { kind: 'error', message: String(P2P_CLOSE_CODES.roomNotFound) }; } if (now - this.lastHostHeartbeatAt > this.limits.hostHeartbeatTimeoutMs) { return { kind: 'error', message: String(P2P_CLOSE_CODES.hostMissing) }; } return null; } private log(event: string, fields: Record): void { this.logger?.info(event, fields); } } function byteLength(value: unknown): number { return new TextEncoder().encode(JSON.stringify(value)).byteLength; } function consumeRate( windows: Map, peerId: string, now: number, windowMs: number, maxCount: number, ): SignalEnvelope | null { const current = windows.get(peerId); if (!current || now - current.startedAt >= windowMs) { windows.set(peerId, { startedAt: now, count: 1 }); return null; } current.count++; if (current.count > maxCount) { return { kind: 'error', message: String(P2P_CLOSE_CODES.rateLimited) }; } return null; } function dayKey(ms: number): string { return new Date(ms).toISOString().slice(0, 10); } function sourceOf(envelope: SignalEnvelope & { target?: string }): string { if ('from' in envelope && typeof envelope.from === 'string') return envelope.from; if ('target' in envelope && typeof envelope.target === 'string') return `peer:${envelope.target}`; return 'anonymous'; }