返回源码地图

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

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

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

1/**
2 * Program evaluation, output capture, and host binding proxies over the process channel.
3 * @module @deepseek-ai/dsh-ptc-runtime-node/src/bootstrap
4 */
5
6import { inspect } from 'node:util'
7import type { DoneMessage, ReplyMessage, ProgramBootData, ProgramToHost } from './protocol.ts'
8import { jsonStringBytesUpTo, jsonValueBytesUpTo, truncateJsonStringBytes } from './output-json.ts'
9import { decodePtcJsonWire, encodePtcJsonWire, snapshotPtcJsonValue } from './json-wire.ts'
10
11const CapturedError = Error
12const capturedObjectCreate = Object.create
13const capturedObjectDefineProperty = Object.defineProperty
14
15/** Define one public binding-error field without consulting mutable globals or descriptor prototypes. */
16function defineBindingErrorField(error: Error, key: string, value: string): void {
17 const attributes = capturedObjectCreate(null) as PropertyDescriptor
18 attributes.enumerable = true
19 attributes.value = value
20 capturedObjectDefineProperty(error, key, attributes)
21}
22
23/** Program messages and host replies transported by the private process channel. */
24export interface BootstrapPort {
25 postMessage(message: ProgramToHost): void
26 on(event: 'message', listener: (message: ReplyMessage) => void): void
27}
28
29/**
30 * A writable stream's `write` slot, as the bootstrap patches it (see
31 * {@link captureStreamWrites}). Method-typed so the real
32 * `process.stdout`/`process.stderr` (narrower chunk parameters) remain
33 * assignable.
34 */
35export interface PatchableStream {
36 write(chunk: unknown, ...rest: unknown[]): boolean
37}
38
39/**
40 * Ordered text capture under the shared outer JSON-byte budget, delivered to
41 * a sink as each item lands (the real sink streams text over the port eagerly,
42 * so captured output survives a mid-run termination). It includes the log
43 * array syntax and string escaping in its accounting. Once exhausted it emits
44 * the fitting prefix and reports the limit once; the host turns that condition
45 * into an explicit `output-limit` run failure.
46 */
47export class LogBuffer {
48 private bytes = 2 // JSON serialization of the empty logs array: []
49 private entries = 0
50 private truncated = false
51 // Explicit fields, not constructor parameter properties: this module loads
52 // under Node's native strip-only mode, which rejects non-erasable syntax —
53 // and parameter properties are non-erasable.
54 private readonly sink: (text: string) => void
55 private readonly onLimit: () => void
56 private readonly maxBytes: number
57
58 constructor(maxBytes: number, sink: (text: string) => void, onLimit: () => void = () => {}) {
59 this.maxBytes = maxBytes
60 this.sink = sink
61 this.onLimit = onLimit
62 }
63
64 /**
65 * Emit text to the sink, charging it against the budget (drops + marks once exhausted).
66 * @param text - the captured text to deliver.
67 */
68 push(text: string): void {
69 if (this.truncated) return
70 const separatorBytes = this.entries > 0 ? 1 : 0
71 const availableBytes = this.maxBytes - this.bytes - separatorBytes
72 const stringBytes = jsonStringBytesUpTo(text, availableBytes)
73 if (stringBytes === undefined) {
74 this.truncated = true
75 const prefix = truncateJsonStringBytes(text, availableBytes)
76 if (prefix.length > 0) {
77 const prefixBytes = jsonStringBytesUpTo(prefix, availableBytes)
78 /* v8 ignore next -- truncateJsonStringBytes guarantees the returned prefix fits. */
79 if (prefixBytes === undefined) throw new CapturedError('program output ledger produced an oversized log prefix')
80 this.bytes += prefixBytes + separatorBytes
81 this.entries += 1
82 this.sink(prefix)
83 }
84 this.onLimit()
85 return
86 }
87 this.bytes += stringBytes + separatorBytes
88 this.entries += 1
89 this.sink(text)
90 }
91
92 /** Remaining exact JSON-byte budget for the completion value or failure message. */
93 remainingOutputBytes(): number {
94 return this.maxBytes - this.bytes
95 }
96}
97
98/** The five console methods the shim captures, in the seam's level vocabulary. */
99const CONSOLE_LEVELS = ['log', 'info', 'warn', 'error', 'debug'] as const
100
101/**
102 * A `console` replacement whose five leveled methods render their arguments
103 * `util.inspect`-style (matching real console formatting closely enough for
104 * a model to recognize its own output) into the buffer. Only these five
105 * exist — the program gets a deliberately small console, not Node's full
106 * console API.
107 * @param logs - the buffer every rendered line is pushed into.
108 * @returns the five-method console object handed to the program.
109 */
110export function makeConsoleShim(logs: LogBuffer): Record<(typeof CONSOLE_LEVELS)[number], (...args: unknown[]) => void> {
111 const render = (args: unknown[]): string =>
112 args.map(arg => typeof arg === 'string' ? arg : inspect(arg, INSPECT_OPTIONS)).join(' ')
113 const shim = Object.create(null) as Record<(typeof CONSOLE_LEVELS)[number], (...args: unknown[]) => void>
114 for (const level of CONSOLE_LEVELS) {
115 shim[level] = (...args: unknown[]) => { logs.push(render(args)) }
116 }
117 return shim
118}
119
120/**
121 * Redirect a stream's `write` into the log buffer (the program-visible
122 * `process.stdout`/`process.stderr` in the Node child), so raw writes land in emission order
123 * alongside console output instead of racing down a pipe. It preserves Node's optional callback
124 * contract: the callback runs asynchronously after admission, even when the log budget drops
125 * the write.
126 *
127 * @param logs - the buffer captured writes are pushed into.
128 * @param stream - the stream whose `write` slot is patched.
129 * @returns the restore function (the in-process tests un-patch; the real
130 * child never needs to).
131 */
132export function captureStreamWrites(logs: LogBuffer, stream: PatchableStream): () => void {
133 // The slot's VALUE is stored for restore and reassigned — never invoked
134 // detached, so the unbound-method concern does not apply.
135 // oxlint-disable-next-line typescript/unbound-method
136 const original = stream.write
137 stream.write = (chunk: unknown, ...rest: unknown[]): boolean => {
138 logs.push(typeof chunk === 'string' ? chunk : String(chunk))
139 // Node's optional-encoding shape: the callback is whichever of the next
140 // two positions holds a function (a non-function there is the encoding).
141 const callback = [rest[0], rest[1]].find(
142 (arg): arg is (error?: Error | null) => void => typeof arg === 'function',
143 )
144 if (callback) queueMicrotask(() => { callback(null) })
145 return true
146 }
147 return () => { stream.write = original }
148}
149
150/** Bounded inspect options: deep enough to be useful, bounded so a pathological value cannot explode the rendering. */
151const INSPECT_OPTIONS = { depth: 4, maxArrayLength: 100, maxStringLength: 10_000 } as const
152
153/**
154 * Prepare the program's completion value for the done message. Only lossless
155 * JSON crosses, and a value that does not fit the remaining combined outer
156 * budget reports `output-limit`; the host revalidates hostile traffic and
157 * remains authoritative for native pipe writes the program shim cannot observe.
158 *
159 * @param value - the program's completion value.
160 * @param remainingOutputBytes - exact bytes left after captured logs.
161 * @param maxOutputBytes - the configured cap named in an overflow diagnostic.
162 * @returns the done-message fragment: `{}` for `undefined`, else a flat wire `{ value }`.
163 */
164export function prepareCompletion(
165 value: unknown,
166 remainingOutputBytes: number,
167 maxOutputBytes: number = remainingOutputBytes,
168): Omit<DoneMessage, 'type'> {
169 if (value === undefined) return {}
170 let snapshot: ReturnType<typeof snapshotPtcJsonValue>
171 try {
172 snapshot = snapshotPtcJsonValue(value)
173 } catch {
174 snapshot = undefined
175 }
176 if (snapshot === undefined) {
177 return prepareFailure(
178 'invalid-output',
179 'program completion must be lossless JSON',
180 remainingOutputBytes,
181 maxOutputBytes,
182 )
183 }
184 if (jsonValueBytesUpTo(snapshot, remainingOutputBytes) === undefined) {
185 return outputLimit(maxOutputBytes)
186 }
187 return { value: encodePtcJsonWire(snapshot) }
188}
189
190/** Build the fixed overflow fragment without carrying rejected variable bytes. */
191function outputLimit(maxOutputBytes: number): Omit<DoneMessage, 'type'> {
192 return { error: { kind: 'output-limit', message: `outer output exceeded ${maxOutputBytes} bytes` } }
193}
194
195/** Admit one bounded failure message or replace it with the fixed overflow diagnostic. */
196function prepareFailure(
197 kind: 'exception' | 'invalid-output',
198 message: string,
199 remainingOutputBytes: number,
200 maxOutputBytes: number,
201): Omit<DoneMessage, 'type'> {
202 if (jsonStringBytesUpTo(message, remainingOutputBytes) === undefined) return outputLimit(maxOutputBytes)
203 return { error: { kind, message } }
204}
205
206/**
207 * Prepare a thrown program value without sending an unbounded stack or
208 * string across the process channel.
209 * @param error - the value thrown by the program.
210 * @param remainingOutputBytes - exact bytes left after captured logs.
211 * @param maxOutputBytes - the configured cap named in an overflow diagnostic.
212 * @returns a bounded exception or fixed output-limit fragment.
213 */
214export function prepareException(
215 error: unknown,
216 remainingOutputBytes: number,
217 maxOutputBytes: number = remainingOutputBytes,
218): Omit<DoneMessage, 'type'> {
219 let message: string
220 try {
221 const detail: unknown = error instanceof CapturedError ? error.stack ?? error.message : error
222 message = typeof detail === 'string' ? detail : String(detail)
223 } catch {
224 message = 'program threw an unrenderable value'
225 }
226 return prepareFailure('exception', message, remainingOutputBytes, maxOutputBytes)
227}
228
229/** One awaited binding call's settlement handles, keyed by call id in the pending map. */
230export interface PendingCall {
231 resolve(value: unknown): void
232 reject(error: Error): void
233}
234
235/** Constructor type for one program-visible binding rejection class. */
236export type BindingErrorConstructor = new (memberName: string, message: string) => Error
237
238/**
239 * Materialize the real error constructor declared by one namespace.
240 * @param descriptor - program-global class name and member-name property.
241 * @returns the constructor injected into the program and used for rejections.
242 */
243function makeBindingErrorClass(
244 descriptor: { name: string; memberNameProperty: string },
245): BindingErrorConstructor {
246 return class BindingCallError extends CapturedError {
247 constructor(memberName: string, message: string) {
248 super(message)
249 defineBindingErrorField(this, 'name', descriptor.name)
250 defineBindingErrorField(this, descriptor.memberNameProperty, memberName)
251 }
252 }
253}
254
255/** Create the namespace-specific rejection for one failed binding call. */
256function bindingFailure(errorClass: BindingErrorConstructor | undefined, memberName: string, message: string): Error {
257 return errorClass ? new errorClass(memberName, message) : new CapturedError(message)
258}
259
260/**
261 * Build each declared error class once so calls and `instanceof` share constructor identity.
262 * @param data - binding namespace declarations from the boot payload.
263 * @returns constructors keyed by their owning namespace global.
264 */
265export function makeBindingErrorClasses(
266 data: Pick<ProgramBootData, 'namespaces'>,
267): Map<string, BindingErrorConstructor> {
268 const classes = new Map<string, BindingErrorConstructor>()
269 for (const namespace of data.namespaces) {
270 if (namespace.errorClass) classes.set(namespace.global, makeBindingErrorClass(namespace.errorClass))
271 }
272 return classes
273}
274
275/**
276 * Route host replies into the pending-call map: each reply settles its call
277 * at most once, and a reply for an unknown id (stray, or a duplicate answer
278 * to an id already settled) is ignored. Shared wiring between
279 * {@link runProgram} and the tests that exercise {@link makeNamespaces}
280 * standalone.
281 * @param port - the port whose `message` events carry the replies.
282 * @param pending - the id-keyed map of unsettled binding calls.
283 */
284export function wireReplies(port: BootstrapPort, pending: Map<number, PendingCall>): void {
285 port.on('message', (message: ReplyMessage) => {
286 const entry = pending.get(message.id)
287 if (!entry) return
288 pending.delete(message.id)
289 if (message.ok) {
290 const value = decodePtcJsonWire(message.value)
291 if (value === undefined) entry.reject(new CapturedError('binding resolution must be lossless JSON'))
292 else entry.resolve(value)
293 } else {
294 entry.reject(new CapturedError(message.message))
295 }
296 })
297}
298
299/**
300 * Build the binding namespace objects the program sees: one null-prototype global per
301 * namespace, each declared name an own enumerable async function that bridges over the port
302 * (`__proto__`/`constructor`/`toString` are ordinary keys, never prototype collisions).
303 * Lossy arguments reject before posting; clone failures and host failure
304 * replies reject only the corresponding call.
305 *
306 * @param data - the boot payload's namespace declarations (globals + names).
307 * @param port - the port binding calls are posted to.
308 * @param pending - the id-keyed map each posted call parks its handles in.
309 * @param nextId - the shared mutable id counter (program-issued correlation ids).
310 * @param errorClasses - per-namespace constructors shared with program globals.
311 * @returns one namespace object per declaration, in declaration order.
312 */
313export function makeNamespaces(
314 data: Pick<ProgramBootData, 'namespaces'>,
315 port: BootstrapPort,
316 pending: Map<number, PendingCall>,
317 nextId: { value: number },
318 errorClasses: Map<string, BindingErrorConstructor> = makeBindingErrorClasses(data),
319): Record<string, unknown>[] {
320 return data.namespaces.map(({ global, names }) => {
321 const errorClass = errorClasses.get(global)
322 const namespace = Object.create(null) as Record<string, unknown>
323 for (const name of names) {
324 Object.defineProperty(namespace, name, {
325 enumerable: true,
326 value: (args: unknown): Promise<unknown> => {
327 let detached: ReturnType<typeof snapshotPtcJsonValue>
328 try {
329 detached = snapshotPtcJsonValue(args)
330 } catch {
331 detached = undefined
332 }
333 if (detached === undefined) {
334 return Promise.reject(bindingFailure(errorClass, name, 'binding arguments must be lossless JSON'))
335 }
336 return new Promise((resolve, reject) => {
337 const id = nextId.value++
338 pending.set(id, {
339 resolve,
340 reject: (error) => {
341 reject(bindingFailure(errorClass, name, error.message))
342 },
343 })
344 try {
345 port.postMessage({ type: 'call', id, global, name, args: encodePtcJsonWire(detached) })
346 } catch (error: unknown) {
347 pending.delete(id)
348 const message = `binding arguments must be structured-cloneable: ${error instanceof CapturedError ? error.message : String(error)}`
349 reject(bindingFailure(errorClass, name, message))
350 }
351 })
352 },
353 })
354 }
355 return namespace
356 })
357}
358
359/**
360 * Run one strict async-function body, allowing top-level `await` and `return`, and post exactly
361 * one terminal {@link DoneMessage}; a thrown program error becomes its `error` field.
362 * @param port - host message port or test double.
363 * @param data - the boot payload the host sent.
364 * @param streams - stdout/stderr objects captured as program logs.
365 * @returns after posting the done message.
366 */
367export async function runProgram(
368 port: BootstrapPort,
369 data: ProgramBootData,
370 streams: { stdout: PatchableStream; stderr: PatchableStream },
371): Promise<void> {
372 const logs = new LogBuffer(
373 data.maxOutputBytes,
374 (text) => { port.postMessage({ type: 'log', text }) },
375 () => { port.postMessage({ type: 'output-limit' }) },
376 )
377 captureStreamWrites(logs, streams.stdout)
378 captureStreamWrites(logs, streams.stderr)
379
380 const pending = new Map<number, PendingCall>()
381 wireReplies(port, pending)
382
383 const nextId = { value: 1 }
384 const errorClasses = makeBindingErrorClasses(data)
385 const namespaces = makeNamespaces(data, port, pending, nextId, errorClasses)
386 const errorClassParameters: string[] = []
387 const errorClassValues: BindingErrorConstructor[] = []
388 for (const namespace of data.namespaces) {
389 if (!namespace.errorClass) continue
390 errorClassParameters.push(namespace.errorClass.name)
391 const errorClass = errorClasses.get(namespace.global)
392 /* v8 ignore next -- makeBindingErrorClasses covers every declaration in the same data. */
393 if (!errorClass) throw new CapturedError(`missing binding error class for ${namespace.global}`)
394 errorClassValues.push(errorClass)
395 }
396 const consoleShim = makeConsoleShim(logs)
397
398 let done: DoneMessage
399 try {
400 // The async function constructor, reached through an instance because
401 // `AsyncFunction` is not a global. The program body is strict-mode.
402 /* v8 ignore next -- the arrow exists only to reach the AsyncFunction constructor; it is never invoked. */
403 const AsyncFunction = (async () => {}).constructor as new (...args: string[]) => (...fnArgs: unknown[]) => Promise<unknown>
404 const fn = new AsyncFunction(
405 ...data.namespaces.map(namespace => namespace.global),
406 ...errorClassParameters,
407 'console',
408 `'use strict';\n${data.code}`,
409 )
410 const value = await fn(...namespaces, ...errorClassValues, consoleShim)
411 done = {
412 type: 'done',
413 ...prepareCompletion(value, logs.remainingOutputBytes(), data.maxOutputBytes),
414 }
415 } catch (error: unknown) {
416 done = {
417 type: 'done',
418 ...prepareException(error, logs.remainingOutputBytes(), data.maxOutputBytes),
419 }
420 }
421 port.postMessage(done)
422}