/** * ProtomuxRPC session wrapping a HyperDHT secret stream. */ import ProtomuxRPC from 'protomux-rpc' import b4a from 'b4a' import { PROTOCOL, PROTOCOL_VERSION, Roles } from '../../shared/protocol.js' import { encodings } from '../../shared/encodings.js' import rateLimiter from '../utils/rateLimiter.js' import logger from '../utils/logger.js' import { resolveRole, assertAllowed, maxRole } from '../core/acl.js' import { audit, shouldAudit } from '../core/audit.js' import { redeemCapability, isPeerAllowed, getPeerEntry, registerPeer, touchPeer, } from '../core/peer-policy.js' import { validateMethodArgs, SCHEMA_VERSION } from '../../shared/schema.js' import { verifyAdminProof } from '../../shared/crypto-auth.js' import { getMacKey, getServerPublicKeyHex } from '../core/auth-keys.js' export class PeerSession { /** * @param {import('stream').Duplex} stream * @param {{ serverPublicKey: Uint8Array, onClose?: (s: PeerSession) => void }} opts */ constructor(stream, { serverPublicKey, onClose } = {}) { this.stream = stream this.id = stream.remotePublicKey ? b4a.toString(stream.remotePublicKey, 'hex') : `anon-${Date.now()}` this.remotePublicKey = stream.remotePublicKey this.closed = false this.onClose = onClose this.role = resolveRole(this.id) this.clientInfo = null this.authMode = 'viewer' this.displayName = null /** @type {Map} */ this.state = new Map() this.rpc = new ProtomuxRPC(stream, { id: serverPublicKey, protocol: PROTOCOL, ...encodings, }) this.rpc.on('close', () => this._handleClose()) this.rpc.on('destroy', () => this._handleClose()) stream.on('close', () => this._handleClose()) stream.on('error', (err) => { logger.error('Peer stream error', { peerId: this.id.slice(0, 12), error: err.message, }) }) } /** * @param {string} method * @param {(args: any, session: PeerSession) => Promise|any} handler * @param {{ hot?: boolean }} [opts] */ respond(method, handler, opts = {}) { const hot = opts.hot === true || rateLimiter.isStreamMethod?.(method) this.rpc.respond(method, encodings, async (args) => { if (!rateLimiter.isAllowed(this, method)) { const err = new Error('Rate limit exceeded. Please wait before making more requests.') err.code = 'RATE_LIMIT_EXCEEDED' throw err } if (hot) { try { assertAllowed(this.role, method) return await handler(args ?? {}, this) } catch (err) { if (err?.code === 'PERMISSION_DENIED') { audit({ method, peerId: this.id, role: this.role, ok: false, error: err.message, force: true, }) } const safe = new Error(sanitizeError(err)) safe.code = err.code || 'UNKNOWN_ERROR' throw safe } } const t0 = Date.now() try { assertAllowed(this.role, method) const validated = validateMethodArgs(method, args ?? {}) if (!validated.ok) { const err = new Error(validated.error || 'Invalid arguments') err.code = 'INVALID_ARGS' throw err } const result = await handler(validated.args, this) if (shouldAudit(method)) { audit({ method, peerId: this.id, role: this.role, ok: true, args: args ?? {}, }) } const latencyMs = Date.now() - t0 if (latencyMs >= 2000) { logger.warn('Slow RPC', { method, peerId: this.id.slice(0, 12), latencyMs, }) } else { logger.debug('RPC ok', { method, peerId: this.id.slice(0, 12), latencyMs, }) } return result } catch (err) { if (shouldAudit(method) || err?.code === 'PERMISSION_DENIED') { audit({ method, peerId: this.id, role: this.role, ok: false, error: err.message, args: args ?? {}, force: err?.code === 'PERMISSION_DENIED', }) } logger.error('RPC handler failed', { method, peerId: this.id.slice(0, 12), role: this.role, error: err.message, code: err.code, }) const safe = new Error(sanitizeError(err)) safe.code = err.code || 'UNKNOWN_ERROR' throw safe } }) } /** * @param {string} method * @param {unknown} payload */ push(method, payload) { if (this.closed || this.rpc.closed) return this.rpc.event(method, payload, encodings) } destroy() { if (this.closed) return this.closed = true try { this.rpc.destroy() } catch { // ignore } try { this.stream.destroy() } catch { // ignore } } _handleClose() { if (this.closed) return this.closed = true if (this.onClose) this.onClose(this) } } /** * @param {PeerSession} session */ export function registerHandshake(session) { session.respond('handshake', async (args) => { if (args?.clientName || args?.clientVersion) { session.clientInfo = { name: args.clientName || 'unknown', version: args.clientVersion || null, } } let authMode = 'viewer' let elevatedRole = null if (args?.adminProof) { const serverPk = getServerPublicKeyHex() const proofRes = verifyAdminProof(getMacKey(), args.adminProof, { peerId: session.id, serverPublicKeyHex: serverPk || '', }) if (!proofRes.ok) { const e = new Error(proofRes.error || 'Admin proof failed') e.code = proofRes.code || 'ADMIN_PROOF_FAILED' audit({ method: 'handshake', peerId: session.id, role: session.role, ok: false, error: e.message, force: true, }) throw e } elevatedRole = Roles.admin authMode = 'seed' registerPeer(session.id, { role: Roles.admin }) } const capabilityToken = args?.capability || null if (capabilityToken && authMode !== 'seed') { try { const { role, reconnected } = redeemCapability(String(capabilityToken), session.id) elevatedRole = role authMode = reconnected ? 'registered' : 'capability' } catch (err) { const registered = getPeerEntry(session.id) if (registered?.role && registered.role !== Roles.viewer) { elevatedRole = registered.role authMode = 'registered' } else { const e = new Error(err.message || 'Capability redeem failed') e.code = err.code || 'CAPABILITY_INVALID' audit({ method: 'handshake', peerId: session.id, role: session.role, ok: false, error: e.message, force: true, }) throw e } } } if (!elevatedRole && authMode === 'viewer') { const registered = getPeerEntry(session.id) if (registered?.role) { elevatedRole = registered.role authMode = 'registered' } } const baseline = resolveRole(session.id) session.role = elevatedRole ? maxRole(baseline, elevatedRole) : baseline session.authMode = authMode if (!isPeerAllowed(session.id, { authMode })) { const err = new Error('Peer not allowed (revoked or not on allowlist)') err.code = 'PEER_DENIED' throw err } touchPeer(session.id) logger.info('Handshake complete', { peerId: session.id.slice(0, 12), role: session.role, authMode, }) audit({ method: 'handshake', peerId: session.id, role: session.role, ok: true, force: true, args: { authMode, role: session.role }, }) return { success: true, protocol: PROTOCOL, protocolVersion: PROTOCOL_VERSION, schemaVersion: SCHEMA_VERSION, role: session.role, peerId: session.id, serverTime: Date.now(), auth: { mode: authMode, role: session.role }, features: { schemaValidation: true, hmacAuth: true, invites: true, room: true, }, } }) } function sanitizeError(err) { const msg = String(err?.message || err || 'Unknown error') return msg.length > 800 ? msg.slice(0, 800) + '…' : msg }