返回源码地图

packages/ptc-runtime/ptc-runtime-node/src/index.ts

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

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

1/** Confined Node programs with host-owned bindings, output limits, and managed process cleanup. */
2import { stripTypeScriptTypes } from 'node:module'
3import { isAbsolute } from 'node:path'
4import { Context } from '@deepseek-ai/cordis'
5import z from '@deepseek-ai/schemastery'
6import { PtcRuntime } from '@deepseek-ai/dsh-ptc-runtime'
7import type { PtcBindingNamespace, PtcJsonValue, PtcRunFailure, PtcRunRequest, PtcRunResult, PtcRunSandbox, PtcRunSpec } from '@deepseek-ai/dsh-ptc-runtime'
8import { MAX_TIMER_DELAY_MS, clampTimeout } from '@deepseek-ai/dsh-timeout'
9import { SandboxUnavailableError, classifyRunnerFailure, isRunnerSpawnFailure } from '@deepseek-ai/dsh-sandbox'
10import type { ConfinedArgv, SandboxExecutionPolicy, SandboxMode } from '@deepseek-ai/dsh-sandbox'
11import type { SubprocessHandle, SubprocessOutcome } from '@deepseek-ai/dsh-subprocess'
12import type {} from '@deepseek-ai/dsh-sandbox-policy'
13import type {} from '@deepseek-ai/dsh-fs'
14import { snapshotJsonValue } from '@deepseek-ai/dsh-util-values'
15import { validateBindings } from './bindings.ts'
16import { JsonChannel } from './channel.ts'
17import { bootstrapArgs } from './launch.ts'
18import type { LaunchConfig } from './launch.ts'
19import { OutputLedger } from './output-ledger.ts'
20import { drainOutput } from './output-stream.ts'
21import { STARTUP_ENVIRONMENT_NAMES } from './environment.ts'
22import { decodePtcJsonWire, encodePtcJsonWire } from './json-wire.ts'
23import type { ProgramBootData } from './protocol.ts'
24
25/** Deployment-varying runtime bounds and launch choices. */
26export interface Config extends LaunchConfig {
27 /** Default elapsed deadline, including nested tool and approval waits. */
28 timeoutMs?: number
29 /** Maximum numeric elapsed budget accepted by resolve. */
30 maxTimeoutMs?: number
31 /** Combined serialized logs, completion and diagnostic byte cap. */
32 maxOutputBytes?: number
33 /** V8 old-generation heap limit in MiB; native allocations are excluded. */
34 maxOldGenerationSizeMb?: number
35 /** Maximum control frame, outstanding argument and queued control-output bytes. */
36 maxMessageBytes?: number
37 /** Maximum simultaneous host binding calls accepted from a program. */
38 maxPendingCalls?: number
39 /** Managed process termination and output-drain grace in milliseconds. */
40 graceMs?: number
41}
42
43type ResolvedConfig = Required<Omit<Config, 'bootstrapPath'>> & Pick<Config, 'bootstrapPath'>
44interface LiveRun { controller: AbortController; finished: Promise<void> }
45const STRIP_PREFIX = 'async function __dsh_program__() {\n'
46const STRIP_SUFFIX = '\n}'
47
48function messageOf(error: unknown): string { return error instanceof Error ? error.message : String(error) }
49function record(value: unknown): value is Record<string, unknown> { return typeof value === 'object' && value !== null && !Array.isArray(value) }
50
51/** Node provider; direct file effects use the same sandbox service as Bash. */
52export class NodePtcRuntime extends PtcRuntime {
53 static inject = ['fs', 'subprocess', 'sandbox', 'sandboxPolicy']
54 static Config: z<Config> = z.object({
55 timeoutMs: z.number().default(120_000),
56 maxTimeoutMs: z.number().default(600_000),
57 maxOutputBytes: z.number().default(67_108_864),
58 maxOldGenerationSizeMb: z.number().default(512),
59 maxMessageBytes: z.number().default(134_217_728),
60 maxPendingCalls: z.number().default(128),
61 graceMs: z.number().default(3_000),
62 nodeExecutable: z.string(),
63 bootstrapPath: z.string(),
64 })
65 readonly language = 'typescript'
66 readonly isolation = 'process'
67 override get executionInstructions(): string {
68 return 'Each call runs in a fresh Node process. Node APIs are available through await import(...). Relative paths use the supplied working directory; process.env starts empty. Direct filesystem access follows this execution\'s sandbox policy.'
69 }
70 private readonly config: ResolvedConfig
71 private readonly live = new Set<LiveRun>()
72 private disposed = false
73
74 constructor(ctx: Context, config: Config) {
75 super(ctx)
76 this.config = { ...config, nodeExecutable: config.nodeExecutable ?? process.execPath } as ResolvedConfig
77 for (const [key, value] of Object.entries(this.config)) {
78 if (typeof value === 'number' && (!Number.isFinite(value) || value <= 0)) throw new Error(`ptc-runtime-node: ${key} must be positive and finite`)
79 }
80 for (const key of ['timeoutMs', 'maxTimeoutMs', 'graceMs'] as const) {
81 if (this.config[key] > MAX_TIMER_DELAY_MS) throw new Error(`ptc-runtime-node: ${key} exceeds the supported timer range`)
82 }
83 if (!Number.isSafeInteger(this.config.maxOutputBytes) || this.config.maxOutputBytes < 4) throw new Error('ptc-runtime-node: maxOutputBytes must be an integer of at least 4')
84 if (!Number.isSafeInteger(this.config.maxMessageBytes) || this.config.maxMessageBytes > 0xffff_ffff) throw new Error('ptc-runtime-node: maxMessageBytes must fit an unsigned 32-bit frame length')
85 if (!Number.isSafeInteger(this.config.maxPendingCalls)) throw new Error('ptc-runtime-node: maxPendingCalls must be an integer')
86 if (!Number.isSafeInteger(this.config.maxOldGenerationSizeMb)) throw new Error('ptc-runtime-node: maxOldGenerationSizeMb must be an integer')
87 if (this.config.nodeExecutable.length === 0) throw new Error('ptc-runtime-node: nodeExecutable must be non-empty')
88 if (this.config.bootstrapPath !== undefined && !isAbsolute(this.config.bootstrapPath)) throw new Error('ptc-runtime-node: bootstrapPath must be absolute')
89 ctx.effect(() => async () => {
90 this.disposed = true
91 const active = [...this.live]
92 for (const run of active) run.controller.abort('runtime disposed')
93 await Promise.all(active.map(run => run.finished))
94 }, 'Node ptc-runtime cleanup')
95 }
96
97 override get sandboxMode(): SandboxMode { return this.ctx.sandboxPolicy.defaultMode }
98
99 override get timeout(): { defaultMs: number; maxMs: number } {
100 return { defaultMs: Math.min(this.config.timeoutMs, this.config.maxTimeoutMs), maxMs: this.config.maxTimeoutMs }
101 }
102
103 /**
104 * Resolve an execution under explicit or deployment policy.
105 * @param request - Program, bindings, optional cwd/deadline, and resolved authority.
106 * @returns Complete execution inputs with a capped numeric budget or an explicit null deadline.
107 */
108 resolve(request: PtcRunRequest): PtcRunSpec {
109 if (this.disposed) throw new Error('ptc-runtime-node: resolve after disposal')
110 const sandboxPolicy = request.sandboxPolicy ?? this.ctx.sandboxPolicy.resolve()
111 const cwd = request.cwd ?? sandboxPolicy.workspaceRoot
112 if (!isAbsolute(cwd)) throw new Error('ptc-runtime-node: cwd must be absolute')
113 return {
114 ...request,
115 cwd,
116 timeoutMs: request.timeoutMs === null ? null : clampTimeout(request.timeoutMs, this.config.timeoutMs, this.config.maxTimeoutMs, 'ptc-runtime-node: timeoutMs'),
117 sandboxPolicy,
118 }
119 }
120
121 /**
122 * Run a resolved program in a fresh managed and confined Node process.
123 * @param spec - Inputs returned by resolve; missing authority is caller misuse.
124 * @returns Output and file-confinement facts after managed cleanup.
125 */
126 async run(spec: PtcRunSpec): Promise<PtcRunResult> {
127 if (this.disposed) throw new Error('ptc-runtime-node: run after disposal')
128 if (spec.sandboxPolicy === undefined) throw new Error('ptc-runtime-node: run requires a resolved sandbox policy')
129 if (!isAbsolute(spec.cwd) || (spec.timeoutMs !== null && (!Number.isFinite(spec.timeoutMs) || spec.timeoutMs <= 0 || spec.timeoutMs > this.config.maxTimeoutMs))) throw new Error('ptc-runtime-node: run requires resolved cwd and timeout')
130 const bindings = validateBindings(spec)
131 const controller = new AbortController()
132 const completion = Promise.withResolvers<void>()
133 const live = { controller, finished: completion.promise }
134 this.live.add(live)
135 try {
136 return await this.execute(spec, spec.sandboxPolicy, bindings, controller)
137 } finally {
138 this.live.delete(live)
139 completion.resolve()
140 }
141 }
142
143 private async execute(
144 spec: PtcRunSpec,
145 policy: SandboxExecutionPolicy,
146 bindings: Map<string, PtcBindingNamespace>,
147 controller: AbortController,
148 ): Promise<PtcRunResult> {
149 const output = new OutputLedger(this.config.maxOutputBytes)
150 const logs: string[] = []
151 const sandbox: PtcRunSandbox = { mode: policy.mode, denied: false }
152 const result = Promise.withResolvers<PtcRunResult>()
153 const signal = spec.signal === undefined ? controller.signal : AbortSignal.any([spec.signal, controller.signal])
154 let handle: SubprocessHandle | undefined
155 let channel: JsonChannel | undefined
156 let confined: ConfinedArgv | undefined
157 let settled = false
158 let timedOut = false
159 let outputOverflow = false
160 let overflowResult: PtcRunResult | undefined
161 let stderr = ''
162 let parsing = true
163 const wallTimer = spec.timeoutMs === null ? undefined
164 : setTimeout(() => { timedOut = true; controller.abort('execution deadline reached') }, spec.timeoutMs)
165 const finish = (failure?: PtcRunFailure, value?: PtcJsonValue): void => {
166 if (settled) return
167 settled = true
168 clearTimeout(wallTimer)
169 signal.removeEventListener('abort', onAbort)
170 channel?.close()
171 void (async () => {
172 if (handle !== undefined) {
173 try {
174 handle.terminate()
175 await Promise.all([handle.done.catch(() => {}), handle.waitForExit()])
176 const drained = await Promise.all([
177 drainOutput(handle.stdout, this.config.graceMs),
178 drainOutput(handle.stderr, this.config.graceMs),
179 ])
180 if (drained.includes(false) && failure === undefined) {
181 failure = { kind: 'worker-exit', message: 'Node process output did not close cleanly' }
182 }
183 } catch (error: unknown) {
184 failure = { kind: 'worker-exit', message: `managed process cleanup failed: ${messageOf(error)}` }
185 } finally {
186 handle.stdout?.destroy()
187 handle.stderr?.destroy()
188 }
189 }
190 const outcome = outputOverflow ? overflowResult ?? output.limit(logs)
191 : failure === undefined ? output.success(logs, value) : output.failure(logs, failure)
192 result.resolve({ ...outcome, sandbox: { ...sandbox } })
193 })()
194 }
195 const onAbort = (): void => {
196 finish(timedOut
197 ? { kind: 'timeout', message: `execution deadline reached (${spec.timeoutMs}ms)` }
198 : { kind: 'abort', message: messageOf(signal.reason) })
199 }
200 signal.addEventListener('abort', onAbort, { once: true })
201 if (signal.aborted) onAbort()
202 try {
203 // Abort callbacks can settle execution before or during an awaited operation.
204 // oxlint-disable-next-line typescript/no-unnecessary-condition
205 if (settled) return await result.promise
206 const stripped = stripTypeScriptTypes(STRIP_PREFIX + spec.program + STRIP_SUFFIX)
207 parsing = false
208 const data: ProgramBootData = {
209 code: stripped.slice(STRIP_PREFIX.length, stripped.length - STRIP_SUFFIX.length),
210 namespaces: [...bindings.values()].map(binding => ({
211 global: binding.global,
212 names: Object.keys(binding.functions),
213 ...binding.errorClass ? { errorClass: binding.errorClass } : {},
214 })),
215 maxOutputBytes: this.config.maxOutputBytes,
216 }
217 const executable = await this.ctx.subprocess.resolveExecutable(this.config.nodeExecutable, undefined, signal)
218 // Abort callbacks can settle execution before or during an awaited operation.
219 // oxlint-disable-next-line typescript/no-unnecessary-condition
220 if (settled) return await result.promise
221 const packaged = 'pkg' in process && this.config.bootstrapPath === undefined
222 const heapFlag = `--max-old-space-size=${this.config.maxOldGenerationSizeMb}`
223 const argv = [executable, ...packaged ? [] : [heapFlag], ...bootstrapArgs(this.ctx.fs, this.config, this.config.maxMessageBytes)]
224 confined = policy.mode === 'danger-full-access' ? undefined : await this.ctx.sandbox.confine(argv, { ...policy, mode: policy.mode }, signal)
225 // oxlint-disable-next-line typescript/no-unnecessary-condition -- Cancellation can settle during awaited confinement.
226 if (settled) return await result.promise
227 if (confined !== undefined) sandbox.enforcement = confined.enforcement
228 // Electron needs its Node-mode selector until bootstrap; the child then removes it with other ambient values.
229 const env: NodeJS.ProcessEnv = Object.fromEntries(Object.keys(process.env)
230 .filter(key => !STARTUP_ENVIRONMENT_NAMES.has(key.toUpperCase()) && key.toUpperCase() !== 'ELECTRON_RUN_AS_NODE')
231 .map(key => [key, undefined]))
232 if (packaged) {
233 env.DSH_PTC_RUNTIME_NODE = '1'
234 env.NODE_OPTIONS = heapFlag
235 }
236 handle = this.ctx.subprocess.spawn({ argv: confined?.argv ?? argv, cwd: spec.cwd, env, stdio: { stdin: 'ignore', stdout: 'pipe', stderr: 'pipe', control: 'pipe' }, graceMs: this.config.graceMs, signal })
237 const launched = handle
238 if (launched.control === undefined || launched.stdout === undefined || launched.stderr === undefined) {
239 throw new Error('subprocess provider did not supply the requested control and output pipes')
240 }
241 const admit = (text: string): void => {
242 if (outputOverflow) return
243 if (!output.admit(text, logs)) {
244 outputOverflow = true
245 overflowResult = output.limit([...logs, text])
246 finish({ kind: 'output-limit', message: `outer output exceeded ${this.config.maxOutputBytes} bytes` })
247 }
248 }
249 const stdoutDecoder = new TextDecoder('utf-8', { ignoreBOM: true })
250 const stderrDecoder = new TextDecoder('utf-8', { ignoreBOM: true })
251 handle.stdout?.on('data', (chunk: Buffer) => {
252 const text = stdoutDecoder.decode(chunk, { stream: true })
253 if (text.length > 0) admit(text)
254 })
255 handle.stdout?.on('end', () => {
256 const text = stdoutDecoder.decode()
257 if (text.length > 0) admit(text)
258 })
259 handle.stdout?.on('error', (error: Error) => { finish({ kind: 'worker-exit', message: messageOf(error) }) })
260 handle.stderr?.on('data', (chunk: Buffer) => {
261 const text = stderrDecoder.decode(chunk, { stream: true })
262 stderr = (stderr + text).slice(-this.config.maxOutputBytes)
263 if (text.length > 0) admit(text)
264 })
265 handle.stderr?.on('end', () => {
266 const text = stderrDecoder.decode()
267 if (text.length > 0) admit(text)
268 })
269 handle.stderr?.on('error', (error: Error) => { finish({ kind: 'worker-exit', message: messageOf(error) }) })
270 let ready = false
271 let nextId = 1
272 let pending = 0
273 let pendingBytes = 0
274 const protocolFailure = (message: string): void => { finish({ kind: 'protocol', message }) }
275 const processFinished = (outcome: SubprocessOutcome): void => {
276 if (settled) return
277 finish((confined !== undefined && classifyRunnerFailure(outcome.exitCode, stderr, confined.runnerFailureRules) !== undefined)
278 ? { kind: 'sandbox-unavailable', message: stderr }
279 : { kind: 'worker-exit', message: `Node process exited before completing (${String(outcome.exitCode)})${stderr ? `: ${stderr}` : ''}` })
280 }
281 const transport: JsonChannel = new JsonChannel(launched.control, this.config.maxMessageBytes, (raw, bytes) => {
282 if (!record(raw)) { protocolFailure('invalid control frame'); return }
283 if (!ready) {
284 if (raw.type !== 'ready') { protocolFailure('program frame arrived before bootstrap readiness'); return }
285 ready = true
286 void transport.send({ type: 'boot', data }).catch((error: unknown) => { protocolFailure(messageOf(error)) })
287 return
288 }
289 switch (raw.type) {
290 case 'log':
291 if (typeof raw.text !== 'string') { protocolFailure('invalid log frame'); return }
292 admit(raw.text)
293 return
294 case 'output-limit': outputOverflow = true; finish(); return
295 case 'done': {
296 if (raw.error !== undefined) {
297 if (!record(raw.error) || typeof raw.error.message !== 'string' || (raw.error.kind !== 'exception' && raw.error.kind !== 'invalid-output' && raw.error.kind !== 'output-limit')) { protocolFailure('invalid terminal error'); return }
298 const failure = { kind: raw.error.kind, message: raw.error.message } as PtcRunFailure
299 if (confined !== undefined) {
300 sandbox.denied = confined.denialSignatures.some(signature =>
301 failure.message.toLowerCase().includes(signature.toLowerCase()))
302 }
303 if (failure.kind === 'output-limit') outputOverflow = true
304 finish(failure)
305 return
306 }
307 const value = raw.value === undefined ? undefined : decodePtcJsonWire(raw.value)
308 if (raw.value !== undefined && value === undefined) finish({ kind: 'invalid-output', message: 'program completion must be lossless JSON' })
309 else finish(undefined, value)
310 return
311 }
312 case 'call': {
313 if (!Number.isSafeInteger(raw.id) || raw.id !== nextId || typeof raw.global !== 'string' || typeof raw.name !== 'string') { protocolFailure('invalid binding call identity'); return }
314 nextId += 1
315 const functions = bindings.get(raw.global)?.functions
316 const fn = functions !== undefined && Object.hasOwn(functions, raw.name) ? functions[raw.name] : undefined
317 if (typeof fn !== 'function') { protocolFailure('program requested an undeclared binding'); return }
318 const args = decodePtcJsonWire(raw.args)
319 if (args === undefined) { protocolFailure('binding arguments must be lossless JSON'); return }
320 if (++pending > this.config.maxPendingCalls || (pendingBytes += bytes) > this.config.maxMessageBytes) { protocolFailure('pending binding calls exceed configured limits'); return }
321 const id = raw.id
322 void (async () => {
323 let reply: unknown
324 try {
325 const value = snapshotJsonValue(await fn(args))
326 if (value === undefined) throw new Error('binding resolution must be lossless JSON')
327 reply = { type: 'reply', id, ok: true, value: encodePtcJsonWire(value) }
328 } catch (error: unknown) {
329 reply = { type: 'reply', id, ok: false, message: messageOf(error) }
330 } finally {
331 pending -= 1
332 pendingBytes -= bytes
333 }
334 if (!settled) await transport.send(reply)
335 })().catch((error: unknown) => { protocolFailure(messageOf(error)) })
336 return
337 }
338 default: protocolFailure('unknown control message')
339 }
340 }, (error, kind) => {
341 if (kind === 'protocol') protocolFailure(messageOf(error))
342 else if (ready) finish({ kind: 'worker-exit', message: messageOf(error) })
343 else void launched.done.then(processFinished, (failure: unknown) => { finish({ kind: 'worker-exit', message: messageOf(failure) }) })
344 })
345 channel = transport
346 void launched.done.then((outcome) => {
347 // Allow queued control-frame callbacks to run before classifying a command exit.
348 setImmediate(() => { processFinished(outcome) })
349 }, (error: unknown) => { finish({ kind: confined !== undefined && isRunnerSpawnFailure(error, confined.argv[0], spec.cwd) ? 'sandbox-unavailable' : 'worker-exit', message: messageOf(error) }) })
350 } catch (error: unknown) {
351 finish({ kind: error instanceof SandboxUnavailableError ? 'sandbox-unavailable' : parsing ? 'exception' : 'worker-exit', message: messageOf(error) })
352 }
353 return await result.promise
354 }
355}
356
357export default NodePtcRuntime