返回源码地图

packages/ssh/subprocess-ssh/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 10;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/** Remote subprocess and PTY handles with independent SSH streams and helper-owned process lifetimes. */
2import { Duplex, PassThrough, type Readable, type Writable } from 'node:stream'
3import type { Socket } from 'node:net'
4import { Context } from '@deepseek-ai/cordis'
5import { SubprocessRuntime, SubprocessExecutableNotFoundError } from '@deepseek-ai/dsh-subprocess'
6import type { SubprocessCollectedOutputs, SubprocessHandle, SubprocessOutcome, SubprocessOutputMode, SubprocessSpawnSpec, SubprocessTerminalHandle, SubprocessTerminalEnvironment, SubprocessTerminalSignal, SubprocessTerminalSpawnSpec } from '@deepseek-ai/dsh-subprocess'
7import { OutputCollector } from '@deepseek-ai/dsh-subprocess-local/output'
8import type { SshConnection } from '@deepseek-ai/dsh-ssh'
9import { doneSchema, foregroundSchema, terminalActivitySchema, outputSnapshotFrameLimit, outputSnapshotSchema, preparedSchema, remotePath, streamEndpointSchema } from '@deepseek-ai/dsh-ssh/schemas'
10import type { SshProcessId } from '@deepseek-ai/dsh-ssh/schemas'
11import { SshRpcPeer, RemoteOperationError } from '@deepseek-ai/dsh-ssh/protocol'
12import { z } from 'zod'
13
14function 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}
17
18class RemoteCleanupError extends AggregateError {}
19
20/** One remote ordinary process; stdin and control remain usable during asynchronous SSH allocation. */
21class RemoteProcess implements SubprocessHandle {
22 readonly stdin: Writable | undefined
23 readonly stdout: Readable | undefined
24 readonly stderr: Readable | undefined
25 readonly control: Duplex | undefined
26 readonly collected: SubprocessCollectedOutputs
27 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 | undefined
37 private sockets: Socket[] = []
38 private quiescent = false
39 private committed = false
40 private termination: Promise<void> | undefined
41 private readonly detachAbort: () => void
42 private spills: { stdout?: string | undefined; stderr?: string | undefined } = {}
43 private readonly updateCollection: Partial<Record<'stdout' | 'stderr', (snapshot: { tail: string; totalBytes: number }, final: boolean) => void>> = {}
44
45 constructor(private readonly ssh: SshConnection, private readonly spec: SubprocessSpawnSpec) {
46 this.stdin = spec.stdio.stdin === 'pipe' ? this.inbound : undefined
47 this.stdout = spec.stdio.stdout === 'pipe' ? this.out : undefined
48 this.stderr = spec.stdio.stderr === 'pipe' ? this.err : undefined
49 const requestedControl = (spec.stdio as typeof spec.stdio & { control?: 'pipe' }).control
50 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 undefined
54 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 = 0
57 let total = 0
58 let finalized = false
59 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) return
63 if (snapshot.totalBytes < total) throw new Error('SSH helper rewound collected output')
64 finalized = final
65 collector = new OutputCollector(mode.maxBytes, name, undefined)
66 collector.push(bytes)
67 base = snapshot.totalBytes - bytes.length
68 total = snapshot.totalBytes
69 }
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.spills
96 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 error
110 })
111 void this.done.catch(() => {})
112 }
113
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.id
120 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 null
132 }))
133 } else {
134 socket.end()
135 const output = name === 'stdout' ? this.out : this.err
136 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 = true
149 } catch (error) {
150 this.terminate()
151 await this.termination?.catch(() => {})
152 for (const socket of this.sockets) socket.destroy()
153 throw error
154 }
155 }
156
157 terminate(): void {
158 if (this.quiescent || this.termination !== undefined) return
159 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 = true
163 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 error
168 })
169 void this.termination.catch(() => {})
170 }
171 }
172
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 }
177
178 async waitForExit(signal?: AbortSignal): Promise<boolean> {
179 if (this.quiescent) return true
180 if (signal?.aborted) return false
181 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 }
188
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.termination
193 this.quiescent = true
194 this.detachAbort()
195 return true
196 }
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 result
201 }
202
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.err
208 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 | undefined
216 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}
227
228/** SSH provider paired with the SSH filesystem; the remote helper selects POSIX process ownership. */
229export 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()
235
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 error
247 })),
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 }
254
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 error
263 }
264 }
265
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 }
273
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 handle
282 }
283
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 }
292
293 private async createTerminal(spec: SubprocessTerminalSpawnSpec, signal: AbortSignal): Promise<SubprocessTerminalHandle> {
294 const ssh = this.ctx.ssh
295 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.id
300 let socket: Socket | undefined
301 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> | undefined
313 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 closing
329 },
330 }
331 signal.addEventListener('abort', abort, { once: true })
332 this.terminals.add(handle)
333 return handle
334 } 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 error
341 }
342 }
343}
344
345export default SshSubprocessRuntime