1
/** Remote subprocess and PTY handles with independent SSH streams and helper-owned process lifetimes. */2
import { Duplex, PassThrough, type Readable, type Writable } from 'node:stream'3
import type { Socket } from 'node:net'4
import { Context } from '@deepseek-ai/cordis'5
import { SubprocessRuntime, SubprocessExecutableNotFoundError } from '@deepseek-ai/dsh-subprocess'6
import type { SubprocessCollectedOutputs, SubprocessHandle, SubprocessOutcome, SubprocessOutputMode, SubprocessSpawnSpec, SubprocessTerminalHandle, SubprocessTerminalEnvironment, SubprocessTerminalSignal, SubprocessTerminalSpawnSpec } from '@deepseek-ai/dsh-subprocess'7
import { OutputCollector } from '@deepseek-ai/dsh-subprocess-local/output'8
import type { SshConnection } from '@deepseek-ai/dsh-ssh'9
import { doneSchema, foregroundSchema, terminalActivitySchema, outputSnapshotFrameLimit, outputSnapshotSchema, preparedSchema, remotePath, streamEndpointSchema } from '@deepseek-ai/dsh-ssh/schemas'10
import type { SshProcessId } from '@deepseek-ai/dsh-ssh/schemas'11
import { SshRpcPeer, RemoteOperationError } from '@deepseek-ai/dsh-ssh/protocol'12
import { z } from 'zod'14
function environment(env?: NodeJS.ProcessEnv): Record<string, string | null> | undefined {15
return env === undefined ? undefined : Object.fromEntries(Object.entries(env).map(([key, value]) => [key, value ?? null]))16
}18
class RemoteCleanupError extends AggregateError {}20
/** One remote ordinary process; stdin and control remain usable during asynchronous SSH allocation. */21
class RemoteProcess implements SubprocessHandle {22
readonly stdin: Writable | undefined23
readonly stdout: Readable | undefined24
readonly stderr: Readable | undefined25
readonly control: Duplex | undefined26
readonly collected: SubprocessCollectedOutputs27
readonly done: Promise<SubprocessOutcome>28
readonly streamsClosed: Promise<void>29
private readonly inbound = new PassThrough()30
private readonly out = new PassThrough()31
private readonly err = new PassThrough()32
private readonly toControl = new PassThrough()33
private readonly fromControl = new PassThrough()34
private readonly controller = new AbortController()35
private readonly started: Promise<void>36
private id: SshProcessId | undefined37
private sockets: Socket[] = []38
private quiescent = false39
private committed = false40
private termination: Promise<void> | undefined41
private readonly detachAbort: () => void42
private spills: { stdout?: string | undefined; stderr?: string | undefined } = {}43
private readonly updateCollection: Partial<Record<'stdout' | 'stderr', (snapshot: { tail: string; totalBytes: number }, final: boolean) => void>> = {}45
constructor(private readonly ssh: SshConnection, private readonly spec: SubprocessSpawnSpec) {46
this.stdin = spec.stdio.stdin === 'pipe' ? this.inbound : undefined47
this.stdout = spec.stdio.stdout === 'pipe' ? this.out : undefined48
this.stderr = spec.stdio.stderr === 'pipe' ? this.err : undefined49
const requestedControl = (spec.stdio as typeof spec.stdio & { control?: 'pipe' }).control50
this.control = requestedControl === undefined ? undefined : Duplex.from({ writable: this.toControl, readable: this.fromControl })51
for (const stream of [this.inbound, this.out, this.err, this.toControl, this.fromControl, this.control]) stream?.on('error', () => {})52
const collect = (name: 'stdout' | 'stderr', stream: PassThrough, mode: SubprocessOutputMode) => {53
if (mode === 'pipe') return undefined54
if (mode === 'inherit') { stream.pipe(name === 'stdout' ? process.stdout : process.stderr, { end: false }); return undefined }55
let collector = new OutputCollector(mode.maxBytes, name, undefined)56
let base = 057
let total = 058
let finalized = false59
this.updateCollection[name] = (snapshot, final): void => {60
const bytes = Buffer.from(snapshot.tail, 'base64')61
if (bytes.length > mode.maxBytes || snapshot.totalBytes < bytes.length) throw new Error('SSH helper returned invalid collected output coordinates')62
if (finalized) return63
if (snapshot.totalBytes < total) throw new Error('SSH helper rewound collected output')64
finalized = final65
collector = new OutputCollector(mode.maxBytes, name, undefined)66
collector.push(bytes)67
base = snapshot.totalBytes - bytes.length68
total = snapshot.totalBytes69
}70
return {71
readFrom: (offset: number) => {72
const snapshot = collector.readFrom(offset - base)73
return {74
...snapshot, nextOffset: snapshot.nextOffset + base,75
...(this.spills[name] === undefined ? {} : { spillPath: this.spills[name] }),76
}77
},78
}79
}80
const stdout = collect('stdout', this.out, spec.stdio.stdout)81
const stderr = collect('stderr', this.err, spec.stdio.stderr)82
this.collected = { ...(stdout === undefined ? {} : { stdout }), ...(stderr === undefined ? {} : { stderr }) }83
const onAbort = (): void => { this.terminate() }84
spec.signal?.addEventListener('abort', onAbort, { once: true })85
this.detachAbort = () => { spec.signal?.removeEventListener('abort', onAbort) }86
this.started = this.start()87
this.streamsClosed = this.started.then(async () => {88
await Promise.all(this.sockets.map(socket => socket.closed ? Promise.resolve() : new Promise<void>((resolve) => {89
socket.once('close', () => { resolve() })90
})))91
}, () => {})92
this.done = this.started.then(async () => {93
const result = await ssh.request('process.done', { id: this.id }, doneSchema, undefined, true)94
await this.drainOutput()95
this.spills = result.spills96
for (const name of ['stdout', 'stderr'] as const) {97
const snapshot = result.collected[name]98
const finish = this.updateCollection[name]99
if ((snapshot === undefined) !== (finish === undefined)) throw new Error('SSH helper returned mismatched output collection modes')100
if (snapshot !== undefined) finish?.(snapshot, true)101
}102
return { exitCode: result.outcome.exitCode, signal: result.outcome.signal as NodeJS.Signals | null }103
}).catch((error: unknown) => {104
this.terminate()105
for (const socket of this.sockets) socket.destroy()106
for (const stream of [this.inbound, this.out, this.err, this.toControl, this.fromControl]) {107
stream.destroy(error instanceof Error ? error : new Error(String(error)))108
}109
throw error110
})111
void this.done.catch(() => {})112
}114
private async start(): Promise<void> {115
this.spec.signal?.throwIfAborted()116
const prepared = await this.ssh.request(117
'process.prepare', { ...this.spec, signal: undefined, env: environment(this.spec.env) }, preparedSchema, this.controller.signal,118
)119
this.id = prepared.id120
try {121
const sockets = await Promise.all(Object.entries(prepared.streams).map(async ([name, path]) =>122
[name, await this.ssh.connectStream(path, this.controller.signal)] as const))123
this.sockets = sockets.map(([, socket]) => socket)124
for (const [name, socket] of sockets) {125
if (name === 'stdout' || name === 'stderr') {126
const mode = this.spec.stdio[name]127
if (typeof mode === 'object') {128
new SshRpcPeer(socket, socket, outputSnapshotFrameLimit(mode.maxBytes), 1, (method, raw) => Promise.resolve().then(() => {129
if (method !== 'snapshot') throw new Error('Unexpected SSH output-stream operation')130
this.updateCollection[name]?.(outputSnapshotSchema.parse(raw), false)131
return null132
}))133
} else {134
socket.end()135
const output = name === 'stdout' ? this.out : this.err136
const closeSocket = (): void => { socket.destroy() }137
output.once('close', closeSocket)138
socket.once('close', () => { output.off('close', closeSocket) })139
if (output.destroyed) closeSocket()140
else socket.pipe(output)141
}142
}143
if (name === 'stdin') this.inbound.pipe(socket)144
if (name === 'control') { this.toControl.pipe(socket); socket.pipe(this.fromControl) }145
}146
this.controller.signal.throwIfAborted()147
await this.ssh.request('process.start', { id: this.id }, z.object({}).strict(), this.controller.signal)148
this.committed = true149
} catch (error) {150
this.terminate()151
await this.termination?.catch(() => {})152
for (const socket of this.sockets) socket.destroy()153
throw error154
}155
}157
terminate(): void {158
if (this.quiescent || this.termination !== undefined) return159
if (!this.committed) this.controller.abort(new Error('SSH process terminated before launch acknowledgement'))160
if (this.id !== undefined) {161
this.termination = this.ssh.request('process.terminate', { id: this.id }, z.null()).then(() => {162
this.quiescent = true163
this.detachAbort()164
}).catch((error: unknown) => {165
// A connection unable to confirm termination must release its helper lease.166
void this.ssh.dispose().catch(() => {})167
throw error168
})169
void this.termination.catch(() => {})170
}171
}173
closeStreams(): void {174
for (const socket of this.sockets) socket.destroy()175
for (const stream of [this.inbound, this.out, this.err, this.toControl, this.fromControl, this.control]) stream?.destroy()176
}178
async waitForExit(signal?: AbortSignal): Promise<boolean> {179
if (this.quiescent) return true180
if (signal?.aborted) return false181
if (signal === undefined) return this.observeExit()182
const cancelled = Promise.withResolvers<boolean>()183
const abort = (): void => { cancelled.resolve(false) }184
signal.addEventListener('abort', abort, { once: true })185
try { return await Promise.race([this.observeExit(signal), cancelled.promise]) }186
finally { signal.removeEventListener('abort', abort) }187
}189
private async observeExit(signal?: AbortSignal): Promise<boolean> {190
try { await this.started } catch {191
// Startup errors remain on done; a prepared process must first finish termination.192
if (this.termination !== undefined) await this.termination193
this.quiescent = true194
this.detachAbort()195
return true196
}197
if (this.termination !== undefined) { await this.termination; return true }198
const result = await this.ssh.request('process.wait', { id: this.id }, z.boolean(), signal, true)199
if (result) { this.quiescent = true; this.detachAbort() }200
return result201
}203
private async drainOutput(): Promise<void> {204
const disposers: Array<() => void> = []205
const outputs = (['stdout', 'stderr'] as const).map((name) => {206
if (typeof this.spec.stdio[name] === 'object') return Promise.resolve()207
const stream = name === 'stdout' ? this.out : this.err208
if (stream.readableEnded || stream.destroyed) return Promise.resolve()209
return new Promise<void>((resolve) => {210
const done = (): void => { resolve() }211
for (const event of ['end', 'close', 'error']) stream.once(event, done)212
disposers.push(() => { for (const event of ['end', 'close', 'error']) stream.off(event, done) })213
})214
})215
let timer: NodeJS.Timeout | undefined216
try {217
await Promise.race([218
Promise.all(outputs),219
new Promise<void>((resolve) => { timer = setTimeout(resolve, this.spec.graceMs) }),220
])221
} finally {222
clearTimeout(timer)223
for (const dispose of disposers) dispose()224
}225
}226
}228
/** SSH provider paired with the SSH filesystem; the remote helper selects POSIX process ownership. */229
export class SshSubprocessRuntime extends SubprocessRuntime {230
static inject = ['ssh']231
private readonly live = new Set<RemoteProcess>()232
private readonly terminals = new Set<SubprocessTerminalHandle>()233
private readonly terminalAllocations = new Set<Promise<SubprocessTerminalHandle>>()234
private readonly lifetime = new AbortController()236
constructor(ctx: Context) {237
super(ctx)238
ctx.effect(() => async () => {239
this.lifetime.abort(new Error('SSH subprocess provider disposed'))240
for (const handle of this.live) handle.terminate()241
const results = await Promise.allSettled([242
...[...this.live].map(async (handle) => {243
try { await handle.waitForExit() } finally { handle.closeStreams() }244
}),245
...[...this.terminalAllocations].map(allocation => allocation.catch((error: unknown) => {246
if (error instanceof RemoteCleanupError) throw error247
})),248
...[...this.terminals].map(handle => handle.terminate()),249
])250
const errors = results.flatMap(result => result.status === 'rejected' ? [result.reason as unknown] : [])251
if (errors.length > 0) throw new AggregateError(errors, 'SSH process cleanup could not be confirmed')252
})253
}255
override async resolveExecutable(command: string, env?: Readonly<Record<string, string>>, signal?: AbortSignal): Promise<string> {256
try {257
return await this.ctx.ssh.request('executable', { command, env }, remotePath, signal)258
} catch (error) {259
if (error instanceof RemoteOperationError && error.code === 'SUBPROCESS_EXECUTABLE_NOT_FOUND') {260
throw new SubprocessExecutableNotFoundError(error.message, { cause: error })261
}262
throw error263
}264
}266
override terminalEnvironment(signal?: AbortSignal): Promise<SubprocessTerminalEnvironment> {267
return this.ctx.ssh.request('terminal.environment', {}, z.object({268
platform: z.enum(['posix', 'windows']), defaultShell: z.string().optional(),269
}).strict().transform(value => ({ platform: value.platform,270
...(value.defaultShell === undefined ? {} : { defaultShell: value.defaultShell }),271
})), signal)272
}274
override spawn(spec: SubprocessSpawnSpec): SubprocessHandle {275
this.lifetime.signal.throwIfAborted()276
spec.signal?.throwIfAborted()277
const handle = new RemoteProcess(this.ctx.ssh, spec)278
this.live.add(handle)279
void handle.done.then(() => handle.waitForExit()).then(() => handle.streamsClosed)280
.then(() => { this.live.delete(handle) }).catch(() => {})281
return handle282
}284
override async spawnTerminal(spec: SubprocessTerminalSpawnSpec): Promise<SubprocessTerminalHandle> {285
this.lifetime.signal.throwIfAborted()286
const signal = spec.signal === undefined ? this.lifetime.signal : AbortSignal.any([spec.signal, this.lifetime.signal])287
signal.throwIfAborted()288
const allocation = this.createTerminal(spec, signal)289
this.terminalAllocations.add(allocation)290
try { return await allocation } finally { this.terminalAllocations.delete(allocation) }291
}293
private async createTerminal(spec: SubprocessTerminalSpawnSpec, signal: AbortSignal): Promise<SubprocessTerminalHandle> {294
const ssh = this.ctx.ssh295
const prepared = await ssh.request('process.prepare', {296
argv: spec.argv, cwd: spec.cwd, env: environment(spec.env), graceMs: spec.graceMs,297
terminal: { rows: spec.rows, cols: spec.cols, terminalType: spec.terminalType, shellActivity: spec.shellActivity },298
}, preparedSchema, signal)299
const id = prepared.id300
let socket: Socket | undefined301
try {302
signal.throwIfAborted()303
socket = await ssh.connectStream(streamEndpointSchema.parse(prepared.streams.terminal), signal)304
signal.throwIfAborted()305
socket.end()306
const output = new PassThrough()307
socket.pipe(output)308
const started = await ssh.request('process.start', { id }, z.object({ pid: z.number().int().positive() }).strict(), signal)309
signal.throwIfAborted()310
const done = ssh.request('process.done', { id }, doneSchema, undefined, true).then(result => ({ exitCode: result.outcome.exitCode, signal: result.outcome.signal as NodeJS.Signals | null }))311
void done.catch(() => {})312
let closing: Promise<void> | undefined313
const abort = (): void => { void handle.terminate().catch(() => { void ssh.dispose().catch(() => {}) }) }314
const handle: SubprocessTerminalHandle = {315
pid: started.pid, output, done,316
resize: async (cols, rows) => { await ssh.request('terminal.resize', { id, cols, rows }, z.null()) },317
write: async (data) => { await ssh.request('terminal.write', { id, value: data }, z.null()) },318
inspectForeground: async () => await ssh.request('terminal.inspect', { id }, foregroundSchema) ?? undefined,319
inspectActivity: () => ssh.request('terminal.activity', { id }, terminalActivitySchema),320
signalForeground: (signal: SubprocessTerminalSignal) => ssh.request('terminal.signal', { id, value: signal }, z.number().int().positive()),321
terminate: () => {322
closing ??= ssh.request('process.terminate', { id }, z.null(), undefined, true).then(() => {323
socket?.destroy()324
output.destroy()325
signal.removeEventListener('abort', abort)326
this.terminals.delete(handle)327
}).catch((error: unknown) => { closing = undefined; throw error })328
return closing329
},330
}331
signal.addEventListener('abort', abort, { once: true })332
this.terminals.add(handle)333
return handle334
} catch (error) {335
socket?.destroy()336
try { await ssh.request('process.terminate', { id }, z.null()) } catch (cleanupError) {337
void ssh.dispose().catch(() => {})338
throw new RemoteCleanupError([error, cleanupError], 'SSH terminal allocation failed and remote cleanup is unknown')339
}340
throw error341
}342
}343
}345
export default SshSubprocessRuntime