Files
peardock/server/rpc/session.js
T

370 lines
11 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 { recordRpc } from '../services/metrics.js'
import {
redeemInvite,
redeemCapability,
isPeerAllowed,
getPeerEntry,
} from '../core/peer-policy.js'
import { validateMethodArgs, SCHEMA_VERSION } from '../../shared/schema.js'
import { sanitizeClientError } from '../utils/dockerErrors.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 {object} opts
* @param {Uint8Array} opts.serverPublicKey
* @param {(session: PeerSession) => void} [opts.onClose]
*/
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
/** @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 {string} method
* @param {(args: any, session: PeerSession) => Promise<any>|any} handler
* @param {{ hot?: boolean }} [opts] - hot path: skip schema/audit/metrics noise (streams)
*/
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'
if (!hot) recordRpc(method, { ok: false, denied: true })
throw err
}
// Hot path: terminal/stream traffic — minimal middleware
if (hot) {
try {
assertAllowed(this.role, method)
const result = await handler(args ?? {}, this)
return result
} 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
recordRpc(method, { ok: true, latencyMs })
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',
})
}
recordRpc(method, {
ok: false,
denied: err?.code === 'PERMISSION_DENIED',
latencyMs: Date.now() - t0,
})
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)
}
}
/**
* Register session handshake (role + protocol version).
* @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
// 1) Admin seed HMAC proof → full admin (never trust client role field)
if (args?.adminProof) {
const serverPk =
getServerPublicKeyHex() ||
(session.stream?.publicKey ? b4a.toString(session.stream.publicKey, 'hex') : null)
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'
}
// 2) HMAC capability grant (pd1 invite embeds this; or direct capability token)
// Persistent by default; reconnect of registered peers never hard-fails on spent jti.
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) {
// Fall back to prior registration (same client identity) so restarts reconnect
const registered = getPeerEntry(session.id)
if (registered?.role && registered.role !== Roles.viewer) {
elevatedRole = registered.role
authMode = 'registered'
logger.info('Capability failed; using registered peer role', {
peerId: session.id.slice(0, 12),
role: registered.role,
code: err.code,
jti: err.jti ? String(err.jti).slice(0, 8) : null,
})
} else {
// Do NOT fall through to bare viewer when a capability was presented and failed —
// that silently grants "success" without the invite role.
const e = new Error(err.message || 'Capability redeem failed')
e.code = err.code || 'CAPABILITY_INVALID'
logger.warn('Capability redeem failed', {
peerId: session.id.slice(0, 12),
code: e.code,
jti: err.jti ? String(err.jti).slice(0, 8) : null,
grantRole: err.grantRole || null,
registered: registered?.role || null,
})
audit({
method: 'handshake',
peerId: session.id,
role: session.role,
ok: false,
error: e.message,
force: true,
})
throw e
}
}
}
// 3) Legacy invite token (capability string or PEARDOCK_LEGACY_INVITES)
if (args?.inviteToken && authMode === 'viewer') {
try {
const entry = redeemInvite(String(args.inviteToken), session.id)
if (entry?.role) {
elevatedRole = entry.role
authMode = String(args.inviteToken).includes('.') ? 'capability' : 'legacy-invite'
}
} catch (err) {
const registered = getPeerEntry(session.id)
if (registered?.role) {
elevatedRole = registered.role
authMode = 'registered'
} else {
const e = new Error(err.message || 'Invite redeem failed')
e.code = err.code || 'INVITE_INVALID'
throw e
}
}
}
// 4) Already-registered peer reconnect without capability (pubkey only)
if (!elevatedRole && authMode === 'viewer') {
const registered = getPeerEntry(session.id)
if (registered?.role) {
elevatedRole = registered.role
authMode = 'registered'
}
}
// Baseline from env / peer policy, then elevate (never demote an elevated grant)
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
}
logger.info('Handshake complete', {
peerId: session.id.slice(0, 12),
role: session.role,
authMode,
elevated: elevatedRole || null,
baseline,
hadCapability: Boolean(capabilityToken),
})
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: {
binaryStreams: true,
schemaValidation: true,
hmacAuth: true,
connectionInvites: true,
},
}
})
}
/**
* Keep operational detail so the UI/job log can tell the user how to fix issues.
* Long Docker messages used to be replaced with a useless generic string.
*/
function sanitizeError(err) {
return sanitizeClientError(err)
}