返回源码地图

packages/ssh/ssh/src/index.ts

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

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

1/** OpenSSH connection owner for one version-matched POSIX helper and its independent forwarded streams. */
2
3import { spawn, execFile, type ChildProcessWithoutNullStreams } from 'node:child_process'
4import { mkdtemp, rm } from 'node:fs/promises'
5import { join } from 'node:path'
6import { createConnection, type Socket } from 'node:net'
7import { Context, Service } from '@deepseek-ai/cordis'
8import schema from '@deepseek-ai/schemastery'
9import { z } from 'zod'
10import { SshRpcPeer, SSH_PROTOCOL_VERSION } from './protocol.ts'
11import { helloSchema, type SshStreamEndpoint } from './schemas.ts'
12import { authenticateStream } from './stream-security.ts'
13
14type Hello = z.infer<typeof helloSchema>
15
16/** Deployment-owned SSH identity and installed helper; no model argument selects these values. */
17export interface Config {
18 /** OpenSSH host alias, including its existing user, key and known-host configuration. */
19 host: string
20 /** Absolute remote Node executable. */
21 node: string
22 /** Absolute path to the installed, bundled helper entry. */
23 helper: string
24 /** SHA-256 of that bundled helper; mismatches refuse the connection. */
25 helperHash: string
26 /** Absolute remote default workspace. */
27 workspace: string
28 /** Optional preinstalled built PTC entry, paired with its expected digest. */
29 bootstrapPath?: string
30 /** SHA-256 of bootstrapPath; both fields must be supplied together. */
31 bootstrapHash?: string
32 /** Connection and administrative-request deadline, at most 2,147,483,647 milliseconds. */
33 requestTimeoutMs?: number
34 /** Maximum JSON payload bytes per helper request or response. */
35 maxFrameBytes?: number
36 /** Maximum ordinary requests; heartbeat and bounded resource cleanup have reserved capacity. */
37 maxPending?: number
38 /** Remote helper lease; loss of heartbeats starts remote managed cleanup. */
39 leaseMs?: number
40}
41
42declare module '@deepseek-ai/cordis' {
43 interface Context { ssh: SshConnection }
44}
45
46/** One non-reconnecting SSH session; loss invalidates all active operations. */
47export class SshConnection extends Service {
48 static Config: schema<Config> = schema.object({
49 host: schema.string().required(), node: schema.string().required(), helper: schema.string().required(),
50 helperHash: schema.string().required(), workspace: schema.string().required(),
51 bootstrapPath: schema.string(), bootstrapHash: schema.string(),
52 requestTimeoutMs: schema.number().default(30_000), maxFrameBytes: schema.number().default(64 * 1024 * 1024),
53 maxPending: schema.number().default(128), leaseMs: schema.number().default(30_000),
54 })
55
56 /** Verified remote helper coordinates; callers must await this before launch. */
57 readonly ready: Promise<Hello>
58 private rpc: SshRpcPeer | undefined
59 private child: ChildProcessWithoutNullStreams | undefined
60 private childClosed: Promise<void> | undefined
61 private directory: string | undefined
62 private heartbeat: NodeJS.Timeout | undefined
63 private closed = false
64 private readonly lifetime = new AbortController()
65 private readonly operations = new Set<Promise<unknown>>()
66 private disposal: Promise<void> | undefined
67 private failure: Error | undefined
68 private sockets = new Set<Socket>()
69 private nextSocket = 0
70 private readonly config: Required<Omit<Config, 'bootstrapPath' | 'bootstrapHash'>> & Pick<Config, 'bootstrapPath' | 'bootstrapHash'>
71 private remote: Hello | undefined
72
73 constructor(ctx: Context, config: Config) {
74 super(ctx, 'ssh')
75 if (process.platform !== 'linux' && process.platform !== 'darwin') throw new Error('SSH runtime requires a POSIX client')
76 this.config = z.object({
77 host: z.string().regex(/^[a-zA-Z0-9][a-zA-Z0-9_.@-]*$/),
78 node: z.string().startsWith('/'), helper: z.string().startsWith('/'), helperHash: z.string().regex(/^[0-9a-f]{64}$/),
79 workspace: z.string().startsWith('/'), requestTimeoutMs: z.number().int().positive().max(2_147_483_647),
80 bootstrapPath: z.string().startsWith('/').optional(), bootstrapHash: z.string().regex(/^[0-9a-f]{64}$/).optional(),
81 maxFrameBytes: z.number().int().positive().max(64 * 1024 * 1024), maxPending: z.number().int().positive().max(128),
82 leaseMs: z.number().int().min(3000).max(600_000),
83 }).refine(value => (value.bootstrapPath === undefined) === (value.bootstrapHash === undefined), 'bootstrapPath and bootstrapHash must be paired')
84 .parse(config) as typeof this.config
85 this.ready = this.start()
86 // Startup uses Node I/O, local validation, and Error-valued RPC failures.
87 void this.ready.catch((error: unknown) => { this.fail(error as Error) })
88 ctx.effect(() => () => this.dispose())
89 }
90
91 /** Hold plugin readiness until the remote identity and helper digest are verified. */
92 async [Service.init](): Promise<void> { await this.ready }
93
94 /** Verified remote Node executable for the paired PTC runtime. */
95 get nodeExecutable(): string {
96 if (this.remote === undefined) throw new Error('SSH helper is not ready')
97 return this.remote.node
98 }
99
100 /** Verified preinstalled PTC entry; unconfigured runtimes fail before program execution. */
101 get bootstrapPath(): string {
102 if (this.remote === undefined || this.config.bootstrapPath === undefined) throw new Error('SSH PTC requires a verified bootstrapPath and bootstrapHash')
103 return this.config.bootstrapPath
104 }
105
106 /**
107 * Send a helper operation; cancellation never replays an ambiguous mutation.
108 * @param method - the private helper operation.
109 * @param params - JSON request fields validated by the helper.
110 * @param result - response validation before returning provider-visible data.
111 * @param signal - cancellation, which does not undo completed remote effects.
112 * @param wait - allow a process observation to outlast the administrative deadline.
113 * @returns the validated remote result.
114 */
115 async request<T>(method: string, params: unknown, result: z.ZodType<T>, signal?: AbortSignal, wait: boolean = false): Promise<T> {
116 this.assertOpen()
117 await this.ready
118 this.assertOpen()
119 const bounded = wait ? signal : signal === undefined
120 ? AbortSignal.timeout(this.config.requestTimeoutMs)
121 : AbortSignal.any([signal, AbortSignal.timeout(this.config.requestTimeoutMs)])
122 return (this.rpc as SshRpcPeer).request(method, params, result, bounded)
123 }
124
125 /**
126 * Forward one authenticated stream through an independent SSH channel.
127 * @param endpoint - private coordinates issued by this connection's helper.
128 * @param signal - cancellation of allocation and the resulting socket.
129 * @returns a paused socket; attach a consumer before resuming it.
130 */
131 async connectStream(endpoint: SshStreamEndpoint, signal?: AbortSignal): Promise<Socket> {
132 return this.track(this.establishStream(endpoint, signal))
133 }
134
135 private async establishStream(endpoint: SshStreamEndpoint, signal?: AbortSignal): Promise<Socket> {
136 const hello = await this.ready
137 this.assertOpen()
138 signal = signal === undefined ? this.lifetime.signal : AbortSignal.any([signal, this.lifetime.signal])
139 const remote = endpoint.path
140 if (!remote.startsWith(`${hello.root}/`) || /[:\r\n\0]/u.test(remote)) throw new Error('SSH helper returned an invalid stream path')
141 signal.throwIfAborted()
142 const local = join(this.directory as string, `s${this.nextSocket++}`)
143 const forward = `${local}:${remote}`
144 const cancelForward = async (): Promise<void> => {
145 // An unavailable master already removed its forwarding listeners.
146 if (!this.closed) await this.controlCommand(['-O', 'cancel', '-L', forward]).catch(() => {})
147 await rm(local, { force: true })
148 }
149 try {
150 await this.controlCommand(['-O', 'forward', '-o', 'ExitOnForwardFailure=yes', '-L', forward], signal)
151 } catch (error) { await cancelForward(); throw error }
152 signal.throwIfAborted()
153 const socket = createConnection({ path: local, allowHalfOpen: true })
154 this.sockets.add(socket)
155 socket.once('close', () => {
156 this.sockets.delete(socket)
157 void this.track(cancelForward()).catch(() => {})
158 })
159 await new Promise<void>((resolve, reject) => {
160 const cleanup = (): void => {
161 signal.removeEventListener('abort', aborted)
162 socket.off('connect', connected)
163 socket.off('error', failed)
164 socket.off('close', closed)
165 }
166 const connected = (): void => { cleanup(); resolve() }
167 const failed = (error: Error): void => { cleanup(); reject(error) }
168 const closed = (): void => { failed(new Error('SSH connection closed before stream establishment')) }
169 const aborted = (): void => { socket.destroy(signal.reason instanceof Error ? signal.reason : new Error(String(signal.reason))) }
170 socket.once('connect', connected)
171 socket.once('error', failed)
172 socket.once('close', closed)
173 signal.addEventListener('abort', aborted, { once: true })
174 })
175 const authenticated = await authenticateStream(socket, endpoint.capability, this.config.requestTimeoutMs, signal)
176 this.sockets.add(authenticated)
177 authenticated.on('error', () => { authenticated.destroy() })
178 authenticated.once('close', () => { this.sockets.delete(authenticated) })
179 return authenticated
180 }
181
182 /** Tear down the helper's remote managed ranges before releasing the SSH master when reachable. */
183 dispose(): Promise<void> {
184 this.disposal ??= this.disposeOnce()
185 return this.disposal
186 }
187
188 private async disposeOnce(): Promise<void> {
189 this.closed = true
190 this.lifetime.abort(new Error('SSH connection is closing'))
191 if (this.heartbeat !== undefined) clearInterval(this.heartbeat)
192 try {
193 await this.ready.catch(() => {})
194 if (this.failure === undefined) await this.rpc?.request('close', {}, z.null(), AbortSignal.timeout(this.config.requestTimeoutMs))
195 } finally {
196 this.rpc?.close()
197 // TLS wrappers release their reads before their underlying sockets close.
198 const socketClosures = [...this.sockets].reverse().map(socket => new Promise<void>((resolve) => {
199 if (socket.closed) resolve()
200 else { socket.once('close', () => { resolve() }); socket.destroy() }
201 }))
202 this.child?.kill('SIGTERM')
203 const force = setTimeout(() => { this.child?.kill('SIGKILL') }, this.config.requestTimeoutMs)
204 try { await this.childClosed } finally { clearTimeout(force) }
205 await Promise.all(socketClosures)
206 while (this.operations.size > 0) await Promise.allSettled([...this.operations])
207 if (this.directory !== undefined) await rm(this.directory, { recursive: true, force: true })
208 }
209 }
210
211 private controlPath(): string { return join(this.directory as string, 'master') }
212
213 private assertOpen(): void {
214 if (this.closed) throw new Error('SSH connection is closed')
215 if (this.failure !== undefined) throw this.failure
216 }
217
218 private track<T>(operation: Promise<T>): Promise<T> {
219 this.operations.add(operation)
220 void operation.finally(() => { this.operations.delete(operation) }).catch(() => {})
221 return operation
222 }
223
224 private async controlCommand(args: string[], signal?: AbortSignal): Promise<void> {
225 const signals = [this.lifetime.signal, AbortSignal.timeout(this.config.requestTimeoutMs)]
226 if (signal !== undefined) signals.push(signal)
227 const combined = AbortSignal.any(signals)
228 combined.throwIfAborted()
229 const result = Promise.withResolvers<undefined>()
230 const command = execFile('ssh', ['-S', this.controlPath(), ...args, this.config.host], {
231 signal: combined, maxBuffer: 64 * 1024,
232 }, (error) => { if (error === null) result.resolve(undefined); else result.reject(error) })
233 const closed = new Promise<void>((resolve) => { command.once('close', () => { resolve() }) })
234 let force: NodeJS.Timeout | undefined
235 const escalate = (): void => {
236 force = setTimeout(() => { command.kill('SIGKILL') }, this.config.requestTimeoutMs)
237 force.unref()
238 }
239 combined.addEventListener('abort', escalate, { once: true })
240 try { await result.promise }
241 finally {
242 await closed
243 combined.removeEventListener('abort', escalate)
244 if (force !== undefined) clearTimeout(force)
245 }
246 }
247
248 private fail(error: Error): void {
249 if (this.failure !== undefined) return
250 this.failure = error
251 this.lifetime.abort(error)
252 if (this.heartbeat !== undefined) clearInterval(this.heartbeat)
253 this.rpc?.close(error)
254 for (const socket of [...this.sockets].reverse()) socket.destroy(error)
255 this.child?.kill('SIGTERM')
256 }
257
258 private async start(): Promise<Hello> {
259 this.directory = await mkdtemp('/tmp/dsh-ssh-')
260 if (this.closed) throw new Error('SSH connection closed before startup')
261 const quote = (value: string): string => `'${value.replaceAll("'", "'\\''")}'`
262 const command = [this.config.node, '--disable-sigusr1', this.config.helper].map(quote).join(' ')
263 const child = spawn('ssh', [
264 '-T', '-M', '-S', this.controlPath(), '-o', 'ControlPersist=no', '-o', 'BatchMode=yes',
265 '-o', 'StrictHostKeyChecking=yes', '-o', 'ForwardAgent=no', '-o', 'ClearAllForwardings=yes',
266 '-o', 'ServerAliveInterval=10', '-o', 'ServerAliveCountMax=3', this.config.host, command,
267 ], { stdio: ['pipe', 'pipe', 'pipe'] })
268 this.child = child
269 this.childClosed = new Promise((resolve) => { child.once('close', () => { resolve() }) })
270 child.stderr.resume() // SSH diagnostics can contain configured paths; operation errors remain structured.
271 child.once('error', (error) => { this.fail(error) })
272 child.once('close', () => { this.fail(new Error('SSH helper disconnected; remote outcomes and cleanup are unknown')) })
273 const rpc = new SshRpcPeer(child.stdout, child.stdin, this.config.maxFrameBytes, this.config.maxPending)
274 this.rpc = rpc
275 rpc.once('closed', (error) => { this.fail(error as Error) })
276 const hello = await rpc.request('hello', {
277 protocol: SSH_PROTOCOL_VERSION, workspace: this.config.workspace, leaseMs: this.config.leaseMs,
278 ...(this.config.bootstrapPath === undefined ? {} : { bootstrapPath: this.config.bootstrapPath }),
279 }, helloSchema, AbortSignal.timeout(this.config.requestTimeoutMs))
280 if (hello.hash !== this.config.helperHash) throw new Error('SSH helper digest differs from the configured artifact')
281 if (hello.bootstrapHash !== this.config.bootstrapHash) throw new Error('SSH PTC bootstrap digest differs from the configured artifact')
282 this.remote = hello
283 let heartbeatPending: Promise<unknown> | undefined
284 this.heartbeat = setInterval(() => {
285 heartbeatPending ??= rpc.request('heartbeat', {}, z.null(), AbortSignal.timeout(this.config.leaseMs / 2))
286 .catch((error: unknown) => { this.fail(error as Error) })
287 .finally(() => { heartbeatPending = undefined })
288 }, Math.floor(this.config.leaseMs / 3))
289 this.heartbeat.unref()
290 return hello
291 }
292}
293
294export default SshConnection