refactor(e2b): compose portable runtime consumers

This commit is contained in:
Tianyi Cui
2026-07-29 02:51:44 +08:00
parent 97496c3d00
commit 2abc9823e7
64 changed files with 1426 additions and 4229 deletions

View File

@@ -8,29 +8,87 @@ import { randomUUID } from 'node:crypto'
import { posix } from 'node:path'
import { Context } from 'cordis'
import { SubprocessService } from '@deepseek-ai/dsh-subprocess'
import type { SubprocessHandle, SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'
import type {
SubprocessHandle,
SubprocessSpawnSpec,
SubprocessTerminalHandle,
SubprocessTerminalSpawnSpec,
} from '@deepseek-ai/dsh-subprocess'
import { quoteE2BShellArg } from '@deepseek-ai/dsh-e2b'
import { E2BSubprocessHandle } from './process.ts'
import { spawnE2BTerminal } from './terminal.ts'
function signalOpts(signal: AbortSignal | undefined): { signal?: AbortSignal } {
return signal === undefined ? {} : { signal }
}
/** E2B command manager registered as `ctx.subprocess`. */
export class E2BSubprocessService extends SubprocessService {
static inject = ['e2b']
private readonly live = new Set<E2BSubprocessHandle>()
private readonly terminals = new Set<SubprocessTerminalHandle>()
/** @inheritdoc */
readonly cwd: string
/** @inheritdoc */
readonly runtimeRoot: string
/** Create the E2B subprocess service and bind its disposal policy. */
constructor(ctx: Context) {
super(ctx)
this.cwd = ctx.e2b.cwd
this.runtimeRoot = ctx.e2b.runtimeRoot
ctx.effect(() => async () => {
const handles = [...this.live]
for (const handle of handles) handle.terminate()
await Promise.all(handles.map(async (handle) => {
await handle.done.catch(() => {})
await handle.waitForExit()
}))
const terminals = [...this.terminals]
const pending: Promise<unknown>[] = []
for (const handle of handles) {
handle.terminate()
pending.push(handle.done.catch(() => {}).then(() => handle.waitForExit()))
}
for (const terminal of terminals) {
terminal.terminate()
pending.push(terminal.waitForExit())
}
this.live.clear()
this.terminals.clear()
await Promise.all(pending)
}, 'e2b subprocess teardown')
}
/** @inheritdoc */
async resolveExecutable(
command: string,
env?: Readonly<Record<string, string>>,
signal?: AbortSignal,
): Promise<string> {
if (command.length === 0) throw new Error('subprocess-e2b: executable name must be non-empty')
signal?.throwIfAborted()
const sandbox = await this.ctx.e2b.getSandbox()
if (posix.isAbsolute(command)) {
await sandbox.commands.run(
`test -f ${quoteE2BShellArg(command)} -a -x ${quoteE2BShellArg(command)}`,
signalOpts(signal),
)
signal?.throwIfAborted()
return command
}
const path = env?.PATH
const prefix = path === undefined ? '' : `PATH=${quoteE2BShellArg(path)} `
const result = await sandbox.commands.run(
`${prefix}command -v -- ${quoteE2BShellArg(command)}`,
signalOpts(signal),
)
signal?.throwIfAborted()
const executable = result.stdout.trim()
if (!posix.isAbsolute(executable) || executable.includes('\n')) {
throw new Error(`subprocess-e2b: executable ${JSON.stringify(command)} did not resolve to one absolute path`)
}
return executable
}
/** @inheritdoc */
spawn(spec: SubprocessSpawnSpec): SubprocessHandle {
const program = spec.argv[0]
@@ -53,6 +111,29 @@ export class E2BSubprocessService extends SubprocessService {
void handle.done.then(release, release).catch(() => {})
return handle
}
/** @inheritdoc */
async spawnTerminal(spec: SubprocessTerminalSpawnSpec): Promise<SubprocessTerminalHandle> {
const program = spec.argv[0]
if (program === undefined || program.length === 0) {
throw new Error('subprocess-e2b: terminal argv must contain a program')
}
for (const [name, value] of [['rows', spec.rows], ['cols', spec.cols], ['graceMs', spec.graceMs]] as const) {
if (!Number.isSafeInteger(value) || value <= 0) {
throw new Error(`subprocess-e2b: terminal ${name} must be a positive safe integer`)
}
}
spec.signal?.throwIfAborted()
const stateDir = posix.join(this.runtimeRoot, 'terminals', randomUUID())
const terminal = await spawnE2BTerminal(this.ctx.e2b, spec, stateDir)
this.terminals.add(terminal)
const release = async (): Promise<void> => {
await terminal.waitForExit()
this.terminals.delete(terminal)
}
void terminal.done.then(release, release).catch(() => {})
return terminal
}
}
export default E2BSubprocessService

View File

@@ -0,0 +1,386 @@
/** E2B PTY allocation and process-session ownership for the subprocess seam. */
import { Buffer } from 'node:buffer'
import { constants } from 'node:os'
import { PassThrough } from 'node:stream'
import { posix } from 'node:path'
import {
CommandExitError,
FileNotFoundError,
quoteE2BShellArg,
} from '@deepseek-ai/dsh-e2b'
import type { CommandHandle, CommandResult, Sandbox } from '@deepseek-ai/dsh-e2b'
import { SENSITIVE_ENV_PATTERN } from '@deepseek-ai/dsh-subprocess'
import type {
SubprocessOutcome,
SubprocessTerminalForeground,
SubprocessTerminalHandle,
SubprocessTerminalSignal,
SubprocessTerminalSpawnSpec,
} from '@deepseek-ai/dsh-subprocess'
import type E2BSandboxService from '@deepseek-ai/dsh-e2b'
const POLL_MS = 20
const TERMINAL_RUNNER_SOURCE = [
'#!/bin/bash',
'set -euo pipefail',
'dsh_state=$1',
'mapfile -d \'\' -t dsh_env < "$dsh_state/environment"',
'mapfile -d \'\' -t dsh_argv < "$dsh_state/argv"',
'rm -f -- "$dsh_state/environment" "$dsh_state/argv" "$dsh_state/runner.bash"',
'if (( ${#dsh_argv[@]} == 0 )); then',
" printf 'terminal runner received empty argv\\n' >&2",
' exit 125',
'fi',
"printf 'ready\\n' > \"$dsh_state/ready\"",
'exec env -i "${dsh_env[@]}" "${dsh_argv[@]}"',
'',
].join('\n')
interface TerminalPaths {
runner: string
environment: string
argv: string
ready: string
}
function signalOpts(signal: AbortSignal | undefined): { signal?: AbortSignal } {
return signal === undefined ? {} : { signal }
}
function delay(ms: number): Promise<void> {
return new Promise(resolve => setTimeout(resolve, ms))
}
function commandSignal(exitCode: number): NodeJS.Signals | null {
const number = exitCode - 128
if (number <= 0) return null
for (const [name, value] of Object.entries(constants.signals)) {
if (value === number) return name as NodeJS.Signals
}
return null
}
function parsePositiveId(value: string, message: string): number {
const raw = value.trim()
const id = Number(raw)
if (!/^[1-9][0-9]*$/.test(raw) || !Number.isSafeInteger(id)) throw new Error(message)
return id
}
function serializeValues(values: readonly string[], kind: string): string {
for (const value of values) {
if (value.includes('\0')) throw new Error(`subprocess-e2b: terminal ${kind} must not contain NUL bytes`)
}
return values.map(value => `${value}\0`).join('')
}
function remoteEnvironment(raw: string, explicit: Readonly<Record<string, string>> | undefined): string {
const environment = new Map<string, string>()
for (const entry of raw.split('\0')) {
if (entry.length === 0) continue
const separator = entry.indexOf('=')
if (separator <= 0) continue
const name = entry.slice(0, separator)
if (name.startsWith('DSH_') || SENSITIVE_ENV_PATTERN.test(name)) continue
environment.set(name, entry.slice(separator + 1))
}
for (const [name, value] of Object.entries(explicit ?? {})) {
if (name.length === 0 || name.includes('=') || name.includes('\0') || value.includes('\0')) {
throw new Error('subprocess-e2b: terminal environment entries require non-empty NUL-free names without = and NUL-free values')
}
environment.set(name, value)
}
return serializeValues([...environment].map(([name, value]) => `${name}=${value}`), 'environment')
}
async function terminalSessionId(sandbox: Sandbox, pid: number, signal?: AbortSignal): Promise<number> {
const result = await sandbox.commands.run(`ps -o sid= -p ${pid}`, signalOpts(signal))
signal?.throwIfAborted()
return parsePositiveId(result.stdout, `subprocess-e2b: cannot resolve process session for terminal ${pid}`)
}
async function waitUntilReady(
sandbox: Sandbox,
paths: TerminalPaths,
completion: Promise<CommandResult>,
signal?: AbortSignal,
): Promise<void> {
const settled = completion.then(() => true, () => true)
for (;;) {
signal?.throwIfAborted()
try {
if ((await sandbox.files.read(paths.ready, signalOpts(signal))).trim() === 'ready') return
} catch (error: unknown) {
if (!(error instanceof FileNotFoundError)) throw error
}
if (await Promise.race([settled, delay(POLL_MS).then(() => false)])) {
throw new Error('subprocess-e2b: terminal exited before publishing readiness')
}
}
}
/** One E2B PTY and all process groups in its remote process session. */
export class E2BTerminalHandle implements SubprocessTerminalHandle {
readonly pid: number
readonly done: Promise<SubprocessOutcome>
private topLevelExited = false
private termination: Promise<void> | undefined
private terminationSignal: NodeJS.Signals | null = null
private removeAbort: (() => void) | undefined
constructor(
private readonly sandbox: Sandbox,
private readonly handle: CommandHandle,
readonly output: PassThrough,
private readonly completion: Promise<CommandResult>,
private readonly sessionId: number,
private readonly stateDir: string,
private readonly graceMs: number,
signal?: AbortSignal,
) {
this.pid = handle.pid
this.done = this.waitForCommand()
void this.done.then(() => { this.terminate() }, () => { this.terminate() })
if (signal !== undefined) {
const onAbort = (): void => { this.terminate() }
signal.addEventListener('abort', onAbort, { once: true })
this.removeAbort = () => { signal.removeEventListener('abort', onAbort) }
if (signal.aborted) this.terminate()
}
}
/** @inheritdoc */
async write(data: Uint8Array): Promise<void> {
if (this.topLevelExited) throw new Error('terminal process has exited')
await this.sandbox.pty.sendInput(this.pid, data)
}
/** @inheritdoc */
async inspectForeground(): Promise<SubprocessTerminalForeground | undefined> {
try {
const result = await this.sandbox.commands.run(`ps -o tpgid= -p ${this.pid}`)
return {
processGroupId: parsePositiveId(
result.stdout,
`subprocess-e2b: cannot resolve foreground process group for terminal ${this.pid}`,
),
// E2B exposes process-table commands but not the /proc memory access
// needed to prove a specific syscall is waiting on fd 0.
inputWaiting: false,
}
} catch (error: unknown) {
if (error instanceof CommandExitError && this.topLevelExited) return undefined
throw error
}
}
/** @inheritdoc */
async signalForeground(signal: SubprocessTerminalSignal): Promise<number> {
const foreground = await this.inspectForeground()
if (foreground === undefined) {
throw new Error(`subprocess-e2b: cannot resolve foreground process group for terminal ${this.pid}`)
}
if (signal === 'SIGKILL' && foreground.processGroupId === this.pid) {
throw new Error('refusing to SIGKILL the terminal shell; terminate the terminal session instead')
}
await this.sandbox.commands.run(`kill -${signal.slice(3)} -- -${foreground.processGroupId}`)
return foreground.processGroupId
}
/** @inheritdoc */
terminate(): void {
this.termination ??= this.closeOnce().catch((error: unknown) => {
this.termination = undefined
throw error
})
void this.termination.catch(() => {})
}
/** @inheritdoc */
async waitForExit(signal?: AbortSignal): Promise<boolean> {
const quiescence = this.termination ?? this.done.then(
() => { this.terminate(); return this.termination },
() => { this.terminate(); return this.termination },
)
if (signal === undefined) {
await quiescence
return true
}
if (signal.aborted) return false
return await new Promise<boolean>((resolve, reject) => {
const onAbort = (): void => { cleanup(); resolve(false) }
const cleanup = (): void => { signal.removeEventListener('abort', onAbort) }
signal.addEventListener('abort', onAbort, { once: true })
void quiescence.then(
() => { cleanup(); resolve(true) },
(error: unknown) => {
cleanup()
reject(error instanceof Error ? error : new Error(String(error)))
},
)
})
}
private async waitForCommand(): Promise<SubprocessOutcome> {
try {
const result = await this.completion
return { exitCode: result.exitCode, signal: null }
} catch (error: unknown) {
if (error instanceof CommandExitError) {
const signal = this.terminationSignal ?? commandSignal(error.exitCode)
return signal === null ? { exitCode: error.exitCode, signal: null } : { exitCode: null, signal }
}
this.output.destroy(error instanceof Error ? error : new Error(String(error)))
throw error
} finally {
this.topLevelExited = true
if (!this.output.destroyed) this.output.end()
}
}
private async sessionProcessGroups(): Promise<number[]> {
const result = await this.sandbox.commands.run(
`ps -eo sid=,pgid= | awk '$1 == ${this.sessionId} { print $2 }'`,
)
const groups = new Set<number>()
for (const raw of result.stdout.trim().split(/\s+/)) {
if (raw.length === 0) continue
const group = parsePositiveId(
raw,
`subprocess-e2b: invalid process group ${JSON.stringify(raw)} in terminal session ${this.sessionId}`,
)
if (group <= 1) {
throw new Error(`subprocess-e2b: unsafe process group ${group} in terminal session ${this.sessionId}`)
}
groups.add(group)
}
return [...groups]
}
private async signalGroups(groups: number[], signal: 'TERM' | 'KILL'): Promise<void> {
if (groups.length === 0) return
try {
await this.sandbox.commands.run(`kill -${signal} -- ${groups.map(group => `-${group}`).join(' ')}`)
} catch (error: unknown) {
if (!(error instanceof CommandExitError)) throw error
}
}
private async awaitSessionEmpty(kill = false): Promise<number[]> {
const deadline = Date.now() + this.graceMs
for (;;) {
const groups = await this.sessionProcessGroups()
if (groups.length === 0 || Date.now() >= deadline) return groups
if (kill) await this.signalGroups(groups, 'KILL')
await delay(Math.min(POLL_MS, Math.max(1, deadline - Date.now())))
}
}
private async closeOnce(): Promise<void> {
let groups = await this.sessionProcessGroups()
if (groups.length > 0) {
this.terminationSignal = 'SIGTERM'
await this.signalGroups(groups, 'TERM')
groups = await this.awaitSessionEmpty()
}
if (groups.length === 0 && !this.topLevelExited) {
await Promise.race([this.done.catch(() => undefined), delay(this.graceMs)])
}
if (groups.length > 0 || !this.topLevelExited) {
this.terminationSignal = 'SIGKILL'
if (!this.topLevelExited) await this.sandbox.pty.kill(this.pid)
groups = await this.awaitSessionEmpty(true)
if (!this.topLevelExited) await Promise.race([this.done.catch(() => undefined), delay(this.graceMs)])
}
if (groups.length > 0) {
throw new Error(`subprocess-e2b: terminal cleanup failed; surviving process groups: ${groups.join(', ')}`)
}
if (!this.topLevelExited) {
throw new Error(`subprocess-e2b: terminal cleanup failed; surviving pid: ${this.pid}`)
}
this.removeAbort?.()
this.removeAbort = undefined
await this.handle.disconnect()
await this.sandbox.files.remove(this.stateDir).catch(() => {})
}
}
/**
* Allocate an E2B PTY, replace its bootstrap shell with the requested argv,
* and return only after the private runner has published readiness.
* @param runtime - Shared E2B sandbox owner.
* @param spec - Fully specified terminal-process request.
* @param stateDir - Private remote directory for one startup transaction.
* @returns The live subprocess terminal handle.
*/
export async function spawnE2BTerminal(
runtime: E2BSandboxService,
spec: SubprocessTerminalSpawnSpec,
stateDir: string,
): Promise<E2BTerminalHandle> {
const sandbox = await runtime.getSandbox()
spec.signal?.throwIfAborted()
const paths: TerminalPaths = {
runner: posix.join(stateDir, 'runner.bash'),
environment: posix.join(stateDir, 'environment'),
argv: posix.join(stateDir, 'argv'),
ready: posix.join(stateDir, 'ready'),
}
const ambient = await sandbox.commands.run('env -0', signalOpts(spec.signal))
const environment = remoteEnvironment(ambient.stdout, spec.env)
const argv = serializeValues(spec.argv, 'argv')
await sandbox.files.makeDir(stateDir)
await sandbox.commands.run(`chmod 700 -- ${quoteE2BShellArg(stateDir)}`, signalOpts(spec.signal))
await sandbox.files.write([
{ path: paths.runner, data: TERMINAL_RUNNER_SOURCE },
{ path: paths.environment, data: environment },
{ path: paths.argv, data: argv },
], signalOpts(spec.signal))
await sandbox.commands.run(
`chmod 600 -- ${quoteE2BShellArg(paths.runner)} ${quoteE2BShellArg(paths.environment)} ${quoteE2BShellArg(paths.argv)}`,
signalOpts(spec.signal),
)
const output = new PassThrough()
let handle: CommandHandle | undefined
let completion: Promise<CommandResult> | undefined
try {
handle = await sandbox.pty.create({
rows: spec.rows,
cols: spec.cols,
cwd: spec.cwd,
envs: { TERM: 'dumb' },
timeoutMs: 0,
...signalOpts(spec.signal),
onData: (data) => { if (!output.destroyed) output.write(Buffer.from(data)) },
})
completion = handle.wait()
void completion.catch(() => {})
if (!Number.isSafeInteger(handle.pid) || handle.pid <= 0) {
throw new Error(`subprocess-e2b: E2B returned invalid terminal pid ${handle.pid}`)
}
const command = `exec /bin/bash ${quoteE2BShellArg(paths.runner)} ${quoteE2BShellArg(stateDir)}\r`
await sandbox.pty.sendInput(handle.pid, Buffer.from(command), signalOpts(spec.signal))
await waitUntilReady(sandbox, paths, completion, spec.signal)
const sessionId = await terminalSessionId(sandbox, handle.pid, spec.signal)
return new E2BTerminalHandle(
sandbox,
handle,
output,
completion,
sessionId,
stateDir,
spec.graceMs,
spec.signal,
)
} catch (error: unknown) {
output.destroy()
if (handle !== undefined) await handle.kill().catch(() => false)
if (completion !== undefined) await completion.catch(() => {})
await sandbox.files.remove(stateDir).catch(() => {})
throw error
}
}