Files
peardata/server/rpc/session.js
T
Raven Scott 015d92a257
Release rolling / release (push) Has been cancelled
CI / test (push) Has been cancelled
first commit
2026-07-18 16:17:38 -04:00

310 lines
8.4 KiB
JavaScript

/**
* 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<string, any>} */
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>|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
}