import { readP2PAccessToken } from '../access-token'; import { isEnvelope, P2P_CLOSE_CODES } from '../protocol'; import { SignalingRelayCoordinator } from './coordinator'; import { DEFAULT_FREE_RELAY_BYTES_PER_ACCOUNT_PER_DAY } from './limits'; import type { SignalEnvelope } from './protocol'; interface CloudflareWebSocket extends WebSocket { accept(): void; } type ErrorSignalEnvelope = Extract; interface DurableObjectStateLike { storage?: { get(key: string): Promise; put(key: string, value: unknown): Promise; }; } interface DurableRoomEnv { VGAI_BILLING?: { fetch(request: Request): Promise }; VGAI_BILLING_SERVICE_TOKEN?: string; P2P_COLYSEUS_RELAY_DISABLED?: string | undefined; P2P_COLYSEUS_RELAY_DAILY_BUDGET_USD?: string | undefined; P2P_COLYSEUS_RELAY_ESTIMATED_USD_PER_GIB?: string | undefined; P2P_COLYSEUS_RELAY_MAX_GLOBAL_BYTES_PER_DAY?: string | undefined; P2P_COLYSEUS_MAX_GLOBAL_REQUESTS_PER_DAY?: string | undefined; P2P_COLYSEUS_MAX_GLOBAL_WEBSOCKET_MESSAGES_PER_DAY?: string | undefined; P2P_COLYSEUS_MAX_ACTIVE_RELAY_SOCKETS?: string | undefined; P2P_COLYSEUS_ACCESS_TOKEN_SECRET?: string | undefined; P2P_COLYSEUS_FREE_RELAY_BYTES_PER_ACCOUNT_PER_DAY?: string | undefined; P2P_COLYSEUS_REQUIRE_ACCOUNT?: string | undefined; } export class P2PColyseusDurableRoom { private readonly coordinator = new SignalingRelayCoordinator(); private readonly sockets = new Map(); /** * peerId → metering principal, learned when a peer authenticates at * host-register / join. The WebSocket relay-data path binds this principal * ONCE at accept time — and under enforcement `acceptWebSocket` REFUSES a * peerId absent from this map, so the bind is always a verified userId, never * an ip-hash fallback. The POST relay-data path uses it only when * account enforcement is OFF (self-host / BYOK); under enforcement that path * ignores this map entirely and meters against the `userId` carried on the * frame's own verified `accessToken` (see `meterPostRelayData`), so a * spoofable per-message `from` is never authoritative for billing. */ private readonly peerPrincipals = new Map(); constructor( private readonly state?: DurableObjectStateLike, private readonly env: DurableRoomEnv = {}, ) {} async fetch(request: Request): Promise { if (isRelayDisabled(this.env)) { return json({ kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }, 503); } const requestBudgetResult = await this.consumeDailyCounter( 'relay-global-requests', 1, getMaxGlobalRequestsPerDay(this.env), ); if (requestBudgetResult) return json(requestBudgetResult, 429); if (request.headers.get('Upgrade')?.toLowerCase() === 'websocket') { return this.acceptWebSocket(request); } if (request.method === 'OPTIONS') { return json(null); } if (request.method === 'GET') { return json(await this.diagnostics()); } if (request.method !== 'POST') { return json({ kind: 'error', message: 'Use POST' }, 405); } let envelope: SignalEnvelope & { from?: string | undefined; ipHash?: string | undefined }; try { envelope = (await request.json()) as SignalEnvelope & { from?: string | undefined; ipHash?: string | undefined; }; } catch { return json({ kind: 'error', message: String(P2P_CLOSE_CODES.badEnvelope) }, 400); } if (!isSignalEnvelope(envelope)) { return json({ kind: 'error', message: String(P2P_CLOSE_CODES.badEnvelope) }, 400); } const ipHash = envelope.ipHash ?? request.headers.get('x-p2p-ip-hash') ?? undefined; return json(await this.handle(ipHash ? { ...envelope, ipHash } : envelope)); } async handle( envelope: SignalEnvelope & { from?: string | undefined; ipHash?: string | undefined }, ): Promise { if (envelope.kind === 'host-register') { const auth = await this.authorize( envelope.accessToken, envelope.roomName, envelope.from, envelope.ipHash, ); if (auth.error) return auth.error; const hostId = envelope.from ?? 'host'; this.rememberPrincipal(hostId, auth.principal); return this.coordinator.registerHost(envelope.roomName, hostId, envelope.ipHash); } if (envelope.kind === 'join-request') { const auth = await this.authorize( envelope.accessToken, envelope.roomName, envelope.from, envelope.ipHash, ); if (auth.error) return auth.error; const clientId = envelope.from ?? 'client'; this.rememberPrincipal(clientId, auth.principal); return this.coordinator.join(envelope.roomName, clientId); } if ( envelope.kind === 'rtc-offer' || envelope.kind === 'rtc-answer' || envelope.kind === 'rtc-ice' || envelope.kind === 'relay-open' || envelope.kind === 'relay-data' ) { return this.forwardSignal(envelope); } if (envelope.kind === 'relay-drain') { return { kind: 'relay-drained', messages: this.coordinator.drain(envelope.peerId) }; } if (envelope.kind === 'heartbeat') { return this.coordinator.heartbeat( envelope.roomId, envelope.from ?? this.coordinator.diagnostics().hostSignalId, ); } return { kind: 'error', message: `Unsupported signal: ${envelope.kind}` }; } private acceptWebSocket(request: Request): Response { if (this.sockets.size >= getMaxActiveRelaySockets(this.env)) { return json({ kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }, 429); } const url = new URL(request.url); const peerId = url.searchParams.get('peerId'); const targetPeerId = url.searchParams.get('targetPeerId'); const roomId = url.searchParams.get('roomId'); if (!peerId || !targetPeerId || !roomId) { return json({ kind: 'error', message: 'Missing peerId, targetPeerId, or roomId' }, 400); } // Under account enforcement the relay-data WebSocket is reachable ONLY by a peer that // already authenticated at host-register/join (both token-gated by authorize()). That // handshake is what binds peerId → verified userId in peerPrincipals. An unauthenticated // peerId has no userId to meter against, so principalForPeer would fall back to ip-hash — // the exact tokenless bypass the POST relay-data path forbids (meterPostRelayData → // authenticatedRelayPrincipal). Refuse the upgrade so both relay-data transports gate // identically. Self-host (REQUIRE_ACCOUNT off) is unchanged: no principal is required. const remembered = this.peerPrincipals.get(peerId); if (getRequireAccount(this.env) && remembered === undefined) { return json({ kind: 'error', message: String(P2P_CLOSE_CODES.unauthorized) }, 401); } const principal = this.principalForPeer( peerId, request.headers.get('x-p2p-ip-hash') ?? undefined, ); const pair = newWebSocketPair(); const client = pair[0]; const server = pair[1]; server.accept(); this.coordinator.ensureRelayPeer(roomId, peerId, targetPeerId); this.coordinator.forward({ kind: 'relay-open', roomId, target: targetPeerId, from: peerId }); this.sockets.set(peerId, server); server.addEventListener('message', async (event) => { if (typeof event.data !== 'string') return; const outcome = await this.meterRelaySocketMessage( principal, roomId, peerId, targetPeerId, event.data, ); if (outcome.kind === 'error') { server.send(JSON.stringify(outcome.error)); server.close(Number(outcome.error.message), outcome.error.message); return; } this.sockets.get(targetPeerId)?.send(JSON.stringify(outcome.payload)); }); server.addEventListener('close', () => { this.sockets.delete(peerId); this.peerPrincipals.delete(peerId); }); return new Response(null, { status: 101, webSocket: client } as ResponseInit); } /** * Meter a single WebSocket relay-data frame: global WS-message count, global * relay bytes, then the SAME per-account daily byte counter the POST relay * path uses. This is the exact body the accepted socket's `message` listener * runs; it is a named method so the WebSocket byte path is unit-testable * without constructing a status-101 `Response` (which the test runtime * rejects). */ private async meterRelaySocketMessage( principal: string | undefined, roomId: string, peerId: string, targetPeerId: string, rawData: string, ): Promise< { kind: 'error'; error: ErrorSignalEnvelope } | { kind: 'forward'; payload: unknown } > { const messageBudget = await this.consumeDailyCounter( 'relay-global-websocket-messages', 1, getMaxGlobalWebSocketMessagesPerDay(this.env), ); if (messageBudget) return { kind: 'error', error: messageBudget }; const inner = JSON.parse(rawData) as unknown; const signalEnvelope: SignalEnvelope & { target: string; from: string } = { kind: 'relay-data', roomId, target: targetPeerId, from: peerId, envelope: inner as never, }; const bytes = byteLength(signalEnvelope); const globalBudget = await this.consumeGlobalRelayBudget(bytes); if (globalBudget) return { kind: 'error', error: globalBudget }; const accountBudget = await this.consumeAccountRelayBytes(principal, bytes); if (accountBudget) return { kind: 'error', error: accountBudget }; const forwardResult = this.coordinator.forward(signalEnvelope); if (forwardResult?.kind === 'error') return { kind: 'error', error: forwardResult }; return { kind: 'forward', payload: inner }; } private async forwardSignal( envelope: SignalEnvelope & { from?: string | undefined; ipHash?: string | undefined }, ): Promise { if (envelope.kind === 'relay-data') { const budgetResult = await this.meterPostRelayData(envelope); if (budgetResult) return budgetResult; } return this.coordinator.forward(envelope) ?? { kind: 'ok' }; } /** * Meter a POST relay-data frame: global relay bytes, then the per-account * daily allotment (the same counter the WebSocket path decrements). * * Under `P2P_COLYSEUS_REQUIRE_ACCOUNT`, the metering principal is the * `userId` recovered from the frame's own verified `accessToken` — NEVER the * client-supplied `envelope.from`, which is spoofable. Because a peer cannot * mint a token bearing another account's `userId`, this makes it IMPOSSIBLE * for peer A's traffic to drain peer B's allotment: `from` is never * authoritative for billing, so spoofing it changes nothing. A missing or * invalid token, or one carrying no `userId`, is refused (`unauthorized`). * * With enforcement OFF (self-host / BYOK) today's allow-when-unset metering * is preserved byte-for-byte: no token is required, and the principal falls * back to the join-time remembered principal (or ip hash) via * {@link principalForPeer}. */ private async meterPostRelayData( envelope: Extract & { ipHash?: string | undefined }, ): Promise { let principal: string | undefined; if (getRequireAccount(this.env)) { const auth = await this.authenticatedRelayPrincipal(envelope.accessToken); if (auth.error) return auth.error; principal = auth.principal; } else { principal = this.principalForPeer(envelope.from, envelope.ipHash); } const bytes = byteLength(envelope); const globalBudget = await this.consumeGlobalRelayBudget(bytes); if (globalBudget) return globalBudget; return this.consumeAccountRelayBytes(principal, bytes); } /** * Resolve the billing principal for a POST relay-data frame under account * enforcement from the frame's OWN token. The token is verified against this * room's bound name and must carry a `userId`; the principal is that `userId` * and nothing else. `envelope.from` is deliberately not consulted — the token * is the only trustworthy sender identity, and binding to it (rather than to * a per-request `from`) is what makes cross-account draining impossible. */ private async authenticatedRelayPrincipal( token: string | undefined, ): Promise<{ error?: ErrorSignalEnvelope; principal?: string }> { const secret = this.env.P2P_COLYSEUS_ACCESS_TOKEN_SECRET; if (!secret || !token) return { error: unauthorized() }; const claims = await readP2PAccessToken(secret, token, { roomName: this.coordinator.boundRoomName, }); if (!claims?.userId) return { error: unauthorized() }; return { principal: claims.userId }; } private async consumeGlobalRelayBudget(bytes: number): Promise { return this.consumeDailyCounter( 'relay-global-bytes', bytes, getDailyGlobalRelayBytes(this.env), ); } private async consumeDailyCounter( prefix: string, amount: number, max: number, ): Promise { if (!Number.isFinite(max) || max <= 0) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } const storage = this.state?.storage; if (!storage) return null; const key = `${prefix}:${dayKey(Date.now())}`; try { const current = (await storage.get(key)) ?? 0; if (current + amount > max) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } await storage.put(key, current + amount); return null; } catch { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } } private async diagnostics() { const diagnostics = this.coordinator.diagnostics(); return { ...diagnostics, quota: { ...diagnostics.quota, maxGlobalRelayBytesPerDay: getDailyGlobalRelayBytes(this.env), maxGlobalRequestsPerDay: getMaxGlobalRequestsPerDay(this.env), maxGlobalWebSocketMessagesPerDay: getMaxGlobalWebSocketMessagesPerDay(this.env), maxActiveRelaySockets: getMaxActiveRelaySockets(this.env), }, dailyUsage: await this.dailyUsage(), }; } private async dailyUsage() { const day = dayKey(Date.now()); const storage = this.state?.storage; if (!storage) { return { day, globalRequests: 0, globalRelayBytes: 0, globalWebSocketMessages: 0, }; } const [globalRequests, globalRelayBytes, globalWebSocketMessages] = await Promise.all([ storage.get(`relay-global-requests:${day}`), storage.get(`relay-global-bytes:${day}`), storage.get(`relay-global-websocket-messages:${day}`), ]); return { day, globalRequests: globalRequests ?? 0, globalRelayBytes: globalRelayBytes ?? 0, globalWebSocketMessages: globalWebSocketMessages ?? 0, }; } /** * Verify a peer's access token when one is PRESENT and resolve its metering * principal. This mirrors {@link meterPostRelayData}'s soft-mode model so that * merely SETTING `P2P_COLYSEUS_ACCESS_TOKEN_SECRET` is non-breaking: a missing * token is ALLOWED while `P2P_COLYSEUS_REQUIRE_ACCOUNT` is off, even with the * secret set (the principal falls back to the ip hash — self-host / BYOK). A * token that IS supplied is still verified against the secret and rejected if * invalid. Only flipping `REQUIRE_ACCOUNT` on makes a token-carried `userId` * mandatory — so setting the secret alone changes nothing until that flip. */ private async authorize( token: string | undefined, roomName: string, peerId: string | undefined, ipHash: string | undefined, ): Promise<{ error?: ErrorSignalEnvelope; principal?: string | undefined }> { const secret = this.env.P2P_COLYSEUS_ACCESS_TOKEN_SECRET; let userId: string | undefined; if (secret && token) { const claims = await readP2PAccessToken(secret, token, { roomName, peerId }); if (!claims) return { error: unauthorized() }; userId = claims.userId; } if (getRequireAccount(this.env) && !userId) return { error: unauthorized() }; return { principal: userId ?? ipHash }; } private rememberPrincipal(peerId: string, principal: string | undefined): void { if (principal !== undefined) this.peerPrincipals.set(peerId, principal); } private principalForPeer( peerId: string | undefined, ipHash: string | undefined, ): string | undefined { const remembered = peerId ? this.peerPrincipals.get(peerId) : undefined; return remembered ?? ipHash; } private async consumeAccountRelayBytes( principal: string | undefined, bytes: number, ): Promise { if (!principal) return null; const cap = await this.consumeDailyCounter( `relay-account-bytes:${principal}`, bytes, getFreeRelayBytesPerAccountPerDay(this.env), ); if (cap || !getRequireAccount(this.env)) return cap; if (!this.env.VGAI_BILLING) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } const serviceToken = this.env.VGAI_BILLING_SERVICE_TOKEN?.trim(); if (!serviceToken) { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } const day = dayKey(Date.now()); const offset = (await this.state?.storage?.get(`relay-account-bytes:${principal}:${day}`)) ?? bytes; const authorization = await this.env.VGAI_BILLING.fetch( new Request('https://billing.internal/usage/authorize', { method: 'POST', headers: billingHeaders(serviceToken), body: JSON.stringify({ userId: principal, kind: 'relay', provider: 'cloudflare', operation: 'relay-bytes', operationId: `relay:${principal}:${day}:${offset}`, providerPrice: { relayBytes: bytes }, }), }), ); const authorizationBody = (await authorization.json().catch(() => undefined)) as | { reservationId?: unknown } | undefined; if (!authorization.ok || typeof authorizationBody?.reservationId !== 'string') { return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } const settled = await this.env.VGAI_BILLING.fetch( new Request('https://billing.internal/usage/settle', { method: 'POST', headers: billingHeaders(serviceToken), body: JSON.stringify({ reservationId: authorizationBody.reservationId, kind: 'relay', provider: 'cloudflare', operation: 'relay-bytes', providerUsage: { relayBytes: bytes }, externalId: `relay:${principal}:${day}:${offset}`, }), }), ); if (settled.ok) return null; await this.env.VGAI_BILLING.fetch( new Request('https://billing.internal/usage/release', { method: 'POST', headers: billingHeaders(serviceToken), body: JSON.stringify({ reservationId: authorizationBody.reservationId, reason: 'relay_settlement_failed', }), }), ); return { kind: 'error', message: String(P2P_CLOSE_CODES.relayQuotaExceeded) }; } } function billingHeaders(token: string): HeadersInit { return { 'Content-Type': 'application/json', 'X-VGAI-Service-Token': token }; } function unauthorized(): ErrorSignalEnvelope { return { kind: 'error', message: String(P2P_CLOSE_CODES.unauthorized) }; } export function getFreeRelayBytesPerAccountPerDay(env: DurableRoomEnv = {}): number { return Math.floor( parsePositiveNumber(env.P2P_COLYSEUS_FREE_RELAY_BYTES_PER_ACCOUNT_PER_DAY) ?? DEFAULT_FREE_RELAY_BYTES_PER_ACCOUNT_PER_DAY, ); } // Accept every plausible truthy spelling (case-insensitive), not just '1'/'true'. This is a // ONE-WAY security flip: if an operator sets it to a value the reader silently ignored (e.g. // `on`, `yes`), enforcement would stay OFF while they believed the room was closed — the worst // possible failure for an access gate. Erring toward "enabled" for any clear affirmative // removes that footgun. Only an explicit falsy/absent value leaves enforcement off. const REQUIRE_ACCOUNT_TRUTHY = new Set(['1', 'true', 'on', 'yes', 'enabled']); export function getRequireAccount(env: DurableRoomEnv = {}): boolean { const value = env.P2P_COLYSEUS_REQUIRE_ACCOUNT; return typeof value === 'string' && REQUIRE_ACCOUNT_TRUTHY.has(value.trim().toLowerCase()); } export function getDailyGlobalRelayBytes(env: DurableRoomEnv = {}): number { const explicit = parsePositiveNumber(env.P2P_COLYSEUS_RELAY_MAX_GLOBAL_BYTES_PER_DAY); if (explicit !== undefined) return Math.floor(explicit); const budgetUsd = parsePositiveNumber(env.P2P_COLYSEUS_RELAY_DAILY_BUDGET_USD) ?? 100; const estimatedUsdPerGib = parsePositiveNumber(env.P2P_COLYSEUS_RELAY_ESTIMATED_USD_PER_GIB) ?? 1; return Math.floor((budgetUsd / estimatedUsdPerGib) * 1024 * 1024 * 1024); } export function getMaxGlobalRequestsPerDay(env: DurableRoomEnv = {}): number { return Math.floor(parsePositiveNumber(env.P2P_COLYSEUS_MAX_GLOBAL_REQUESTS_PER_DAY) ?? 250_000); } export function getMaxGlobalWebSocketMessagesPerDay(env: DurableRoomEnv = {}): number { return Math.floor( parsePositiveNumber(env.P2P_COLYSEUS_MAX_GLOBAL_WEBSOCKET_MESSAGES_PER_DAY) ?? 2_000_000, ); } export function getMaxActiveRelaySockets(env: DurableRoomEnv = {}): number { return Math.floor(parsePositiveNumber(env.P2P_COLYSEUS_MAX_ACTIVE_RELAY_SOCKETS) ?? 200); } function json(value: unknown, status = 200): Response { return new Response(JSON.stringify(value), { status, headers: { 'Content-Type': 'application/json', 'Access-Control-Allow-Origin': '*', 'Access-Control-Allow-Methods': 'GET,POST,OPTIONS', 'Access-Control-Allow-Headers': 'Content-Type', }, }); } function isSignalEnvelope(value: unknown): value is SignalEnvelope { if (!isRecord(value) || typeof value['kind'] !== 'string') return false; switch (value['kind']) { case 'host-register': case 'join-request': return typeof value['roomName'] === 'string' && value['roomName'].length > 0; case 'rtc-offer': case 'rtc-answer': return ( typeof value['roomId'] === 'string' && typeof value['target'] === 'string' && isRecord(value['sdp']) ); case 'rtc-ice': return ( typeof value['roomId'] === 'string' && typeof value['target'] === 'string' && isRecord(value['candidate']) ); case 'relay-open': return typeof value['roomId'] === 'string' && typeof value['target'] === 'string'; case 'relay-data': return ( typeof value['roomId'] === 'string' && typeof value['target'] === 'string' && isEnvelope(value['envelope']) ); case 'relay-drain': return typeof value['peerId'] === 'string'; case 'heartbeat': return typeof value['roomId'] === 'string'; default: return false; } } function isRecord(value: unknown): value is Record { return typeof value === 'object' && value !== null && !Array.isArray(value); } function newWebSocketPair(): [CloudflareWebSocket, CloudflareWebSocket] { const Pair = ( globalThis as unknown as { WebSocketPair?: new () => Record<'0' | '1', CloudflareWebSocket> } ).WebSocketPair; if (!Pair) throw new Error('WebSocketPair is not available in this runtime'); const pair = new Pair(); return [pair[0], pair[1]]; } function parsePositiveNumber(value: string | undefined): number | undefined { if (value === undefined || value.trim() === '') return undefined; const parsed = Number(value); return Number.isFinite(parsed) && parsed > 0 ? parsed : undefined; } function isRelayDisabled(env: DurableRoomEnv): boolean { return env.P2P_COLYSEUS_RELAY_DISABLED === 'true' || env.P2P_COLYSEUS_RELAY_DISABLED === '1'; } function dayKey(time: number): string { return new Date(time).toISOString().slice(0, 10); } function byteLength(value: unknown): number { return new TextEncoder().encode(JSON.stringify(value)).byteLength; }