1
/** OpenSSH connection owner for one version-matched POSIX helper and its independent forwarded streams. */3
import { spawn, execFile, type ChildProcessWithoutNullStreams } from 'node:child_process'4
import { mkdtemp, rm } from 'node:fs/promises'5
import { join } from 'node:path'6
import { createConnection, type Socket } from 'node:net'7
import { Context, Service } from '@deepseek-ai/cordis'8
import schema from '@deepseek-ai/schemastery'9
import { z } from 'zod'10
import { SshRpcPeer, SSH_PROTOCOL_VERSION } from './protocol.ts'11
import { helloSchema, type SshStreamEndpoint } from './schemas.ts'12
import { authenticateStream } from './stream-security.ts'14
type Hello = z.infer<typeof helloSchema>16
/** Deployment-owned SSH identity and installed helper; no model argument selects these values. */17
export interface Config {18
/** OpenSSH host alias, including its existing user, key and known-host configuration. */19
host: string20
/** Absolute remote Node executable. */21
node: string22
/** Absolute path to the installed, bundled helper entry. */23
helper: string24
/** SHA-256 of that bundled helper; mismatches refuse the connection. */25
helperHash: string26
/** Absolute remote default workspace. */27
workspace: string28
/** Optional preinstalled built PTC entry, paired with its expected digest. */29
bootstrapPath?: string30
/** SHA-256 of bootstrapPath; both fields must be supplied together. */31
bootstrapHash?: string32
/** Connection and administrative-request deadline, at most 2,147,483,647 milliseconds. */33
requestTimeoutMs?: number34
/** Maximum JSON payload bytes per helper request or response. */35
maxFrameBytes?: number36
/** Maximum ordinary requests; heartbeat and bounded resource cleanup have reserved capacity. */37
maxPending?: number38
/** Remote helper lease; loss of heartbeats starts remote managed cleanup. */39
leaseMs?: number40
}42
declare module '@deepseek-ai/cordis' {43
interface Context { ssh: SshConnection }44
}46
/** One non-reconnecting SSH session; loss invalidates all active operations. */47
export 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
})56
/** Verified remote helper coordinates; callers must await this before launch. */57
readonly ready: Promise<Hello>58
private rpc: SshRpcPeer | undefined59
private child: ChildProcessWithoutNullStreams | undefined60
private childClosed: Promise<void> | undefined61
private directory: string | undefined62
private heartbeat: NodeJS.Timeout | undefined63
private closed = false64
private readonly lifetime = new AbortController()65
private readonly operations = new Set<Promise<unknown>>()66
private disposal: Promise<void> | undefined67
private failure: Error | undefined68
private sockets = new Set<Socket>()69
private nextSocket = 070
private readonly config: Required<Omit<Config, 'bootstrapPath' | 'bootstrapHash'>> & Pick<Config, 'bootstrapPath' | 'bootstrapHash'>71
private remote: Hello | undefined73
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.config85
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
}91
/** Hold plugin readiness until the remote identity and helper digest are verified. */92
async [Service.init](): Promise<void> { await this.ready }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.node98
}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.bootstrapPath104
}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.ready118
this.assertOpen()119
const bounded = wait ? signal : signal === undefined120
? 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
}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
}135
private async establishStream(endpoint: SshStreamEndpoint, signal?: AbortSignal): Promise<Socket> {136
const hello = await this.ready137
this.assertOpen()138
signal = signal === undefined ? this.lifetime.signal : AbortSignal.any([signal, this.lifetime.signal])139
const remote = endpoint.path140
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 authenticated180
}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.disposal186
}188
private async disposeOnce(): Promise<void> {189
this.closed = true190
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
}211
private controlPath(): string { return join(this.directory as string, 'master') }213
private assertOpen(): void {214
if (this.closed) throw new Error('SSH connection is closed')215
if (this.failure !== undefined) throw this.failure216
}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 operation222
}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 | undefined235
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 closed243
combined.removeEventListener('abort', escalate)244
if (force !== undefined) clearTimeout(force)245
}246
}248
private fail(error: Error): void {249
if (this.failure !== undefined) return250
this.failure = error251
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
}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 = child269
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 = rpc275
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 = hello283
let heartbeatPending: Promise<unknown> | undefined284
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 hello291
}292
}294
export default SshConnection