Files
peardock/server/handlers/terminal.js
T
snxraven 25ba70cce8 Optimize terminal input: coalesce keystrokes and hot-path RPC
Batch typing into short frames, send UTF-8 instead of base64 for text,
exempt stream methods from the general rate limit, and skip heavy
middleware on terminalInput/resize for lower latency.
2026-07-10 21:39:55 -04:00

264 lines
7.2 KiB
JavaScript

/**
* Interactive container terminal over protomux-rpc.
* TTY sessions stream raw PTY bytes (no demux). Supports multi-session via sessionId.
*/
import { PassThrough } from 'stream'
import { docker } from '../services/docker.js'
import { Pushes } from '../../shared/protocol.js'
import logger from '../utils/logger.js'
const SESSIONS_KEY = 'terminals'
function getSessions(session) {
let map = session.state.get(SESSIONS_KEY)
if (!map) {
map = new Map()
session.state.set(SESSIONS_KEY, map)
}
return map
}
function resolveSessionId(args) {
return args.sessionId || args.terminalId || args.containerId || 'default'
}
/**
* @param {Buffer|Uint8Array} chunk
*/
function toBase64(chunk) {
return Buffer.isBuffer(chunk) ? chunk.toString('base64') : Buffer.from(chunk).toString('base64')
}
export function registerTerminalHandlers(session) {
session.respond('startTerminal', async (args) => {
const containerId = args.containerId
if (!containerId) throw new Error('containerId required')
const sessions = getSessions(session)
const sessionId = resolveSessionId(args)
if (sessions.has(sessionId)) {
endOne(sessions, sessionId)
}
const useTty = args.tty !== false
const shellCandidates = []
if (Array.isArray(args.cmd) && args.cmd.length) {
shellCandidates.push(args.cmd.map(String))
} else if (typeof args.cmd === 'string' && args.cmd.trim()) {
shellCandidates.push(args.cmd.trim().split(/\s+/))
} else if (args.shell) {
shellCandidates.push([String(args.shell)])
} else {
shellCandidates.push(['/bin/bash'], ['/bin/sh'], ['/bin/ash'])
}
const container = docker.getContainer(containerId)
let exec = null
let lastErr = null
for (const Cmd of shellCandidates) {
try {
exec = await container.exec({
Cmd,
AttachStdin: true,
AttachStdout: true,
AttachStderr: true,
Tty: useTty,
})
break
} catch (err) {
lastErr = err
}
}
if (!exec) {
throw lastErr || new Error('Failed to create exec session')
}
const stream = await exec.start({
hijack: true,
stdin: true,
Tty: useTty,
})
const entry = {
containerId,
exec,
stream,
sessionId,
tty: useTty,
}
sessions.set(sessionId, entry)
session.state.set('terminal', entry)
const pushOut = (chunk, isErr = false) => {
const channel = isErr ? Pushes.terminalErrorOutput : Pushes.terminalOutput
const type = isErr ? 'terminalErrorOutput' : 'terminalOutput'
try {
session.push(channel, {
type,
containerId,
sessionId,
data: toBase64(chunk),
encoding: 'base64',
})
} catch (e) {
logger.debug('terminal push failed', { error: e.message })
}
}
if (useTty) {
// Raw PTY — do NOT demux (would corrupt ANSI / binary)
stream.on('data', (chunk) => pushOut(chunk, false))
} else {
const stdout = new PassThrough()
const stderr = new PassThrough()
container.modem.demuxStream(stream, stdout, stderr)
stdout.on('data', (chunk) => pushOut(chunk, false))
stderr.on('data', (chunk) => pushOut(chunk, true))
}
stream.on('end', () => {
sessions.delete(sessionId)
if (session.state.get('terminal') === entry) session.state.delete('terminal')
})
stream.on('error', (err) => {
logger.error('Terminal stream error', { containerId, error: err.message })
sessions.delete(sessionId)
})
const cols = Number(args.cols)
const rows = Number(args.rows)
if (cols > 0 && rows > 0) {
try {
await exec.resize({ h: rows, w: cols })
} catch {
// ignore
}
}
logger.info('Terminal session started', {
containerId,
sessionId,
tty: useTty,
peer: session.id.slice(0, 12),
})
return {
success: true,
message: `Terminal started for ${containerId}`,
sessionId,
containerId,
tty: useTty,
}
})
// Hot path: events (id=0) — keep handler sync/cheap; no audit/metrics in session layer
session.respond(
'terminalInput',
(args) => {
const sessions = getSessions(session)
const sessionId = resolveSessionId(args)
const entry = sessions.get(sessionId) || session.state.get('terminal')
if (!entry) return null
if (args.containerId && args.containerId !== entry.containerId) return null
let inputData
if (args.encoding === 'base64') {
inputData = Buffer.from(args.data || '', 'base64')
} else {
// Default utf8 string in JSON (optimal for keystrokes)
inputData = Buffer.from(args.data || '', 'utf8')
}
if (inputData.length && entry.stream && !entry.stream.writableEnded) {
entry.stream.write(inputData)
}
return null
},
{ hot: true }
)
session.respond(
'terminalResize',
async (args) => {
const sessions = getSessions(session)
const sessionId = resolveSessionId(args)
const entry = sessions.get(sessionId) || session.state.get('terminal')
if (!entry) return null
if (args.containerId && args.containerId !== entry.containerId) return null
const cols = Number(args.cols)
const rows = Number(args.rows)
if (cols > 1 && rows > 0) {
try {
await entry.exec.resize({ h: rows, w: cols })
} catch {
// ignore transient resize errors
}
}
return null
},
{ hot: true }
)
session.respond('killTerminal', async (args) => {
const sessions = getSessions(session)
const sessionId = args.sessionId || args.terminalId
if (sessionId && sessions.has(sessionId)) {
const entry = sessions.get(sessionId)
endOne(sessions, sessionId)
return {
success: true,
message: `Terminal session ${sessionId} killed`,
sessionId,
containerId: entry.containerId,
}
}
if (args.containerId) {
let killed = 0
for (const [id, entry] of [...sessions.entries()]) {
if (entry.containerId === args.containerId) {
endOne(sessions, id)
killed += 1
}
}
if (killed) {
return { success: true, message: `Killed ${killed} terminal(s) for ${args.containerId}` }
}
}
const entry = session.state.get('terminal')
if (entry) {
endOne(sessions, entry.sessionId || 'default')
session.state.delete('terminal')
return { success: true, message: 'Terminal killed' }
}
return { success: false, message: 'No terminal session found' }
})
}
function endOne(sessions, sessionId) {
const entry = sessions.get(sessionId)
if (!entry) return
try {
entry.stream.end()
} catch {
// ignore
}
try {
entry.stream.destroy?.()
} catch {
// ignore
}
sessions.delete(sessionId)
}
export function endTerminal(session) {
const sessions = session.state.get(SESSIONS_KEY)
if (sessions) {
for (const id of [...sessions.keys()]) endOne(sessions, id)
}
session.state.delete('terminal')
session.state.delete(SESSIONS_KEY)
}
export function cleanupTerminalOnClose(session) {
endTerminal(session)
}