forked from snxraven/peardock
370 lines
11 KiB
JavaScript
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)
|
|
}
|