1
0

index.ts 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478
  1. /**
  2. * Worker-thread code runtime: a fresh worker runs each host-type-stripped TypeScript program
  3. * and bridges bindings over its message port. This is containment, not a security boundary:
  4. * model code has bash-equivalent trust despite an empty environment, a heap cap, measured
  5. * event-loop busy-time and wall-time budgets, and termination that also stops synchronous loops.
  6. * @module @deepseek-ai/dsh-code-runtime-worker
  7. */
  8. import { Worker } from 'node:worker_threads'
  9. import { stripTypeScriptTypes } from 'node:module'
  10. import { fileURLToPath } from 'node:url'
  11. import { Context } from 'cordis'
  12. import z from 'schemastery'
  13. import { CodeRuntime } from '@deepseek-ai/dsh-code-runtime'
  14. import type { CodeBindingFunction, CodeJsonValue, CodeRunFailure, CodeRunRequest, CodeRunResult } from '@deepseek-ai/dsh-code-runtime'
  15. import { snapshotJsonValue } from '@deepseek-ai/dsh-session'
  16. import type { ReplyMessage, WorkerBootData, WorkerToHost } from './protocol.ts'
  17. import { truncateJsonStringBytes } from './output-json.ts'
  18. /** Plugin config: every execution cap, changeable from `cordis.yml` (no hardcoded tunables). */
  19. export interface Config {
  20. /**
  21. * Busy-time budget in milliseconds: the run fails with kind `'timeout'`
  22. * once the worker's MEASURED event-loop active time
  23. * (`worker.performance.eventLoopUtilization()`) exceeds this. Metering
  24. * measured busy time — not wall time, not host-side pending-call
  25. * bookkeeping — is what makes the budget both fair (a program awaiting a
  26. * slow tool accrues nothing) and ungameable (a hot loop accrues whether
  27. * or not a decoy dispatch is in flight).
  28. */
  29. computeMs?: number
  30. /**
  31. * Wall-clock ceiling in milliseconds; never pauses for anything. The
  32. * backstop for what busy-time cannot see (a program awaiting a promise
  33. * nobody will resolve).
  34. */
  35. maxWallMs?: number
  36. /** Hard cap for the combined serialized outer logs, completion value, and failure diagnostic. */
  37. maxOutputBytes?: number
  38. /** The worker's max old-generation heap in MiB (`resourceLimits`); overflow kills the worker, surfacing as kind `'worker-exit'`. */
  39. maxOldGenerationSizeMb?: number
  40. }
  41. /** {@link Config} after schemastery fills the defaults (every field present). */
  42. type ResolvedConfig = Required<Config>
  43. /**
  44. * How often the host samples the worker's event-loop utilization for the
  45. * `computeMs` budget. An internal cadence, not config: the only effect of
  46. * the interval is budget-expiry granularity (a run can overshoot by up to
  47. * one interval), and nothing a deployment could tune here improves that
  48. * without burning host CPU.
  49. */
  50. const ELU_POLL_INTERVAL_MS = 25
  51. /** Smallest cap that can represent the empty logs array plus an empty JSON failure diagnostic. */
  52. const MIN_OUTPUT_BYTES = 4
  53. /** ECMAScript reserved words that cannot be async-function parameter names — rejected as binding globals. */
  54. const RESERVED_WORDS = new Set([
  55. 'await', 'break', 'case', 'catch', 'class', 'const', 'continue', 'debugger', 'default', 'delete', 'do',
  56. 'else', 'enum', 'export', 'extends', 'false', 'finally', 'for', 'function', 'if', 'import', 'in',
  57. 'instanceof', 'new', 'null', 'return', 'super', 'switch', 'this', 'throw', 'true', 'try', 'typeof',
  58. 'var', 'void', 'while', 'with', 'yield', 'let', 'static', 'implements', 'interface', 'package',
  59. 'private', 'protected', 'public', 'arguments', 'eval',
  60. ])
  61. /** Valid async-function parameter name (the binding global becomes one). */
  62. const IDENTIFIER = /^[A-Za-z_$][A-Za-z0-9_$]*$/
  63. /**
  64. * The shell a program is wrapped in for the type-strip, matching the
  65. * grammatical context it will execute in (an async function body, where
  66. * top-level `return` and `await` are legal — a bare module parse would
  67. * reject the `return`). Strip mode is position-preserving (removed syntax
  68. * becomes whitespace, nothing shifts), so the wrapper survives the strip
  69. * byte-identical and the body slices back out with the model's own
  70. * line/column positions intact.
  71. */
  72. const STRIP_WRAP = { prefix: 'async function __dsh_program__() {\n', suffix: '\n}' } as const
  73. /** One in-flight run's host-side state, tracked for disposal. */
  74. interface LiveRun {
  75. worker: Worker
  76. settle(failure: CodeRunFailure): void
  77. finished: Promise<void>
  78. }
  79. /**
  80. * The worker entry path. Source runs unbuilt (`src/worker.ts`, loadable
  81. * directly on this repo's Node range via native type stripping — the file
  82. * is erasable-only with type-only relative imports); the built package
  83. * ships it as a sibling CommonJS bundle (`lib/worker.cjs`, its own tsdown
  84. * entry) because pkg's VFS Worker hook compiles string-path entries as
  85. * CommonJS.
  86. * The URL *pathname*'s extension says which world this module is in —
  87. * pathname, because dev-time module runners (vitest) may suffix
  88. * `import.meta.url` with a query string; relative resolution drops it. Worker
  89. * receives a filesystem string so pkg's VFS Worker hook can resolve it.
  90. */
  91. /* v8 ignore next -- the './worker.cjs' arm is the built-lib world, unreachable unbuilt by construction; the built-lib e2e pins it. */
  92. const WORKER_PATH = fileURLToPath(new URL(new URL(import.meta.url).pathname.endsWith('.ts') ? './worker.ts' : './worker.cjs', import.meta.url))
  93. /** Render an unknown thrown value as a message, `Error` or not. */
  94. function messageOf(error: unknown): string {
  95. return error instanceof Error ? error.message : String(error)
  96. }
  97. /**
  98. * Runtime shape gate for inbound port traffic. The peer runs MODEL CODE and
  99. * can post anything — `null`, primitives, objects with poisoned fields — so
  100. * the compile-time `WorkerToHost` type means nothing here: everything is
  101. * re-validated and REBUILT field by field (a forged extra field never rides
  102. * along; a non-number call id can never be echoed into a reply). Junk returns
  103. * `undefined` and is dropped — a throw in the host's `message` listener would
  104. * crash the host process.
  105. */
  106. function parseWorkerMessage(raw: unknown): WorkerToHost | undefined {
  107. if (typeof raw !== 'object' || raw === null) return undefined
  108. const m = raw as Record<string, unknown>
  109. switch (m.type) {
  110. case 'call': {
  111. if (typeof m.id !== 'number' || typeof m.global !== 'string' || typeof m.name !== 'string') return undefined
  112. return { type: 'call', id: m.id, global: m.global, name: m.name, args: m.args }
  113. }
  114. case 'log': {
  115. if (typeof m.text !== 'string') return undefined
  116. return { type: 'log', text: m.text }
  117. }
  118. case 'output-limit': return { type: 'output-limit' }
  119. case 'done': {
  120. if (m.error === undefined) return { type: 'done', ...m.value !== undefined ? { value: m.value } : {} }
  121. const error = m.error
  122. if (typeof error !== 'object' || error === null) return undefined
  123. const { kind, message } = error as Record<string, unknown>
  124. if ((kind !== 'exception' && kind !== 'invalid-output' && kind !== 'output-limit') || typeof message !== 'string') return undefined
  125. return { type: 'done', error: { kind, message } }
  126. }
  127. default: return undefined
  128. }
  129. }
  130. /** Serialized byte size of one lossless JSON value. */
  131. function jsonBytes(value: CodeJsonValue): number {
  132. return Buffer.byteLength(JSON.stringify(value), 'utf8')
  133. }
  134. /** One run's combined outer-output ledger; binding values never enter it. */
  135. class OutputLedger {
  136. private bytes = 2 // JSON serialization of the empty logs array: []
  137. private entries = 0
  138. constructor(private readonly maxBytes: number) {}
  139. /** Admit one exact log entry, or report that the hard cap was crossed. */
  140. admit(text: string, sink: string[]): boolean {
  141. const cost = Buffer.byteLength(JSON.stringify(text), 'utf8') + (this.entries > 0 ? 1 : 0)
  142. if (this.bytes + cost > this.maxBytes) return false
  143. this.bytes += cost
  144. this.entries += 1
  145. sink.push(text)
  146. return true
  147. }
  148. /** Finalize a successful absent-or-JSON completion against the combined cap. */
  149. success(logs: string[], value?: CodeJsonValue): CodeRunResult {
  150. if (value !== undefined && this.bytes + jsonBytes(value) > this.maxBytes) return this.limit(logs)
  151. return { logs, ...value !== undefined ? { value } : {} }
  152. }
  153. /** Finalize a failure diagnostic, with output-limit taking precedence when combined bytes exceed the cap. */
  154. failure(logs: string[], error: CodeRunFailure): CodeRunResult {
  155. if (this.bytes + Buffer.byteLength(JSON.stringify(error.message), 'utf8') > this.maxBytes) return this.limit(logs)
  156. return { logs, error }
  157. }
  158. /** Build the explicit output-limit failure while retaining a fitting prefix of the final log. */
  159. limit(logs: string[]): CodeRunResult {
  160. const fullMessage = `outer output exceeded ${this.maxBytes} bytes`
  161. const messageBytes = Buffer.byteLength(JSON.stringify(fullMessage), 'utf8')
  162. const retained = [...logs]
  163. let retainedBytes = jsonBytes(retained)
  164. const logBudget = this.maxBytes - messageBytes
  165. while (retained.length > 0 && retainedBytes > logBudget) {
  166. const removed = retained.pop()
  167. /* v8 ignore next -- the while guard proves pop cannot return undefined. */
  168. if (removed === undefined) throw new Error('output ledger lost its final log entry')
  169. const separatorBytes = retained.length > 0 ? 1 : 0
  170. retainedBytes -= Buffer.byteLength(JSON.stringify(removed), 'utf8') + separatorBytes
  171. const prefix = truncateJsonStringBytes(removed, logBudget - retainedBytes - separatorBytes)
  172. if (prefix.length > 0) {
  173. retained.push(prefix)
  174. retainedBytes += Buffer.byteLength(JSON.stringify(prefix), 'utf8') + separatorBytes
  175. break
  176. }
  177. }
  178. if (logBudget < 2) {
  179. retained.length = 0
  180. retainedBytes = 2
  181. }
  182. const availableMessageBytes = this.maxBytes - retainedBytes
  183. // This fixed diagnostic is ASCII with no JSON escapes, so two bytes are
  184. // the surrounding quotes and every retained character costs one byte.
  185. const message = messageBytes <= availableMessageBytes
  186. ? fullMessage
  187. : fullMessage.slice(0, availableMessageBytes - 2)
  188. return { logs: retained, error: { kind: 'output-limit', message } }
  189. }
  190. }
  191. /**
  192. * The shipped {@link CodeRuntime} backend (`ctx.codeRuntime`). Registers as
  193. * the `codeRuntime` service; every cap comes from validated config. See the
  194. * module doc for the containment model and the class JSDoc on the seam for
  195. * the contract this implements (error-as-field, hostile-peer port,
  196. * no cross-run state, dispose to quiescence).
  197. */
  198. export class WorkerCodeRuntime extends CodeRuntime {
  199. static Config: z<Config> = z.object({
  200. computeMs: z.number().default(60_000),
  201. maxWallMs: z.number().default(600_000),
  202. maxOutputBytes: z.number().default(67_108_864),
  203. maxOldGenerationSizeMb: z.number().default(512),
  204. })
  205. readonly language = 'typescript'
  206. readonly isolation = 'worker-thread'
  207. private readonly config: ResolvedConfig
  208. private readonly live = new Set<LiveRun>()
  209. private disposed = false
  210. constructor(ctx: Context, config: Config) {
  211. super(ctx)
  212. // Schemastery filled the defaults; the cast records that. Positivity is a
  213. // semantic check the schema's plain number type does not carry.
  214. this.config = config as ResolvedConfig
  215. for (const [key, value] of Object.entries(this.config)) {
  216. if (!(Number.isFinite(value) && value > 0)) throw new Error(`dsh-code-runtime-worker: config.${key} must be a positive number, got ${String(value)}`)
  217. }
  218. if (!Number.isSafeInteger(this.config.maxOutputBytes) || this.config.maxOutputBytes < MIN_OUTPUT_BYTES) {
  219. throw new Error(`dsh-code-runtime-worker: config.maxOutputBytes must be a safe integer of at least ${MIN_OUTPUT_BYTES}, got ${String(this.config.maxOutputBytes)}`)
  220. }
  221. ctx.effect(() => () => this.teardown(), 'worker code-runtime teardown')
  222. }
  223. /**
  224. * Dispose to quiescence: mark the service unusable, fail every in-flight
  225. * run as aborted, and AWAIT each worker's exit so no worker outlives the
  226. * fiber.
  227. */
  228. private async teardown(): Promise<void> {
  229. this.disposed = true
  230. const runs = [...this.live]
  231. for (const run of runs) run.settle({ kind: 'abort', message: 'runtime disposed' })
  232. await Promise.all(runs.map(run => run.finished))
  233. }
  234. /**
  235. * Execute one program in a fresh worker. Program outcomes — including a
  236. * type-strip syntax error, which never spawns a worker — resolve with
  237. * `result.error`; the method rejects only for seam misuse (a disposed
  238. * runtime, an invalid binding namespace).
  239. * @param request - the program, its bindings, and the abort signal.
  240. * @returns the run's outcome per the seam contract.
  241. */
  242. async run(request: CodeRunRequest): Promise<CodeRunResult> {
  243. if (this.disposed) throw new Error('dsh-code-runtime-worker: run() after disposal')
  244. const bindings = this.validateBindings(request)
  245. if (request.signal?.aborted) {
  246. return { logs: [], error: { kind: 'abort', message: String(request.signal.reason) } }
  247. }
  248. let code: string
  249. try {
  250. const stripped = stripTypeScriptTypes(STRIP_WRAP.prefix + request.program + STRIP_WRAP.suffix)
  251. code = stripped.slice(STRIP_WRAP.prefix.length, stripped.length - STRIP_WRAP.suffix.length)
  252. } catch (error: unknown) {
  253. // A program that does not survive the type-strip (syntax error,
  254. // non-erasable syntax like `enum`) is a program failure, reported the
  255. // same way a thrown exception would be — and no worker ever spawns.
  256. return { logs: [], error: { kind: 'exception', message: messageOf(error) } }
  257. }
  258. return await this.execute(request, code, bindings)
  259. }
  260. /** Reject (seam misuse) malformed binding namespaces: non-identifier or reserved globals, duplicates, and the `console` collision. */
  261. private validateBindings(request: CodeRunRequest): Map<string, Record<string, CodeBindingFunction>> {
  262. const bindings = new Map<string, Record<string, CodeBindingFunction>>()
  263. for (const namespace of request.bindings) {
  264. if (!IDENTIFIER.test(namespace.global) || RESERVED_WORDS.has(namespace.global)) {
  265. throw new Error(`dsh-code-runtime-worker: binding global ${JSON.stringify(namespace.global)} is not a usable identifier`)
  266. }
  267. if (namespace.global === 'console' || namespace.global === 'ToolCallError' || bindings.has(namespace.global)) {
  268. throw new Error(`dsh-code-runtime-worker: duplicate binding global ${JSON.stringify(namespace.global)}`)
  269. }
  270. bindings.set(namespace.global, namespace.functions)
  271. }
  272. return bindings
  273. }
  274. /** Spawn the worker for one validated, type-stripped run and drive it to settlement. */
  275. private execute(
  276. request: CodeRunRequest,
  277. code: string,
  278. bindings: Map<string, Record<string, CodeBindingFunction>>,
  279. ): Promise<CodeRunResult> {
  280. const bootData: WorkerBootData = {
  281. code,
  282. namespaces: [...bindings].map(([global, functions]) => ({ global, names: Object.keys(functions) })),
  283. maxOutputBytes: this.config.maxOutputBytes,
  284. }
  285. const worker = new Worker(WORKER_PATH, {
  286. workerData: bootData,
  287. // Model code gets NO ambient environment — stronger than the scrubbed
  288. // env the defensive-patterns rule requires for spawned commands.
  289. env: {},
  290. // Hermetic flags too: without this the worker inherits the host process's execArgv (a
  291. // test runner's or tsx's loader hooks), which a bare isolate with an empty environment
  292. // cannot satisfy.
  293. execArgv: [],
  294. resourceLimits: { maxOldGenerationSizeMb: this.config.maxOldGenerationSizeMb },
  295. // Backstop capture: the bootstrap patches JS-level writes into its own
  296. // ordered buffer, so these pipes normally stay silent; anything that
  297. // still arrives (native-level writes) is appended after the done logs.
  298. stdout: true,
  299. stderr: true,
  300. })
  301. return new Promise<CodeRunResult>((resolve) => {
  302. let settled = false
  303. const answered = new Set<number>()
  304. const logs: string[] = []
  305. const strayLogs: string[] = []
  306. const output = new OutputLedger(this.config.maxOutputBytes)
  307. // No settled guard: `finish` snapshots the arrays when it resolves, so
  308. // a chunk flushing after settlement mutates only the discarded buffers,
  309. // and the ledger bounds that growth until the pipes close.
  310. const captureStray = (chunk: Buffer): void => {
  311. const text = chunk.toString('utf8')
  312. if (!settled && !output.admit(text, strayLogs)) finish(output.limit([...logs, ...strayLogs, text]))
  313. }
  314. worker.stdout.on('data', captureStray)
  315. worker.stderr.on('data', captureStray)
  316. // Exactly one outcome wins. Every path cleans up, terminates, and awaits the worker;
  317. // logs captured before timeout, abort, or failure remain in the result.
  318. let finishResolve!: () => void
  319. const finished = new Promise<void>((done) => { finishResolve = done })
  320. const finish = (result: CodeRunResult): void => {
  321. if (settled) return
  322. settled = true
  323. clearInterval(eluTimer)
  324. clearTimeout(wallTimer)
  325. request.signal?.removeEventListener('abort', onAbort)
  326. this.live.delete(live)
  327. void worker.terminate().then(() => {
  328. finishResolve()
  329. resolve(result)
  330. })
  331. }
  332. const onDone = (message: WorkerToHost): void => {
  333. if (message.type !== 'done') return
  334. const captured = [...logs, ...strayLogs]
  335. if (message.error) {
  336. finish(output.failure(captured, message.error))
  337. return
  338. }
  339. if (message.value === undefined) {
  340. finish(output.success(captured))
  341. return
  342. }
  343. // The worker-thread boundary has already structured-cloned this
  344. // hostile value, so accessors and proxies cannot survive to throw
  345. // during the lossless-JSON snapshot.
  346. const value = snapshotJsonValue(message.value) as CodeJsonValue | undefined
  347. finish(value === undefined
  348. ? output.failure(captured, { kind: 'invalid-output', message: 'program completion must be lossless JSON' })
  349. : output.success(captured, value))
  350. }
  351. const onCall = (message: WorkerToHost): void => {
  352. if (message.type !== 'call' || settled) return
  353. // Hostile-peer rules: a duplicate id is ignored, an unknown name is
  354. // answered with a failure, and a binding throw/reject becomes the
  355. // program-side rejection — contained here, never a host crash.
  356. if (answered.has(message.id)) return
  357. answered.add(message.id)
  358. const reply = (payload: ReplyMessage): void => {
  359. if (settled) return
  360. // Canonical resolutions were snapshotted as lossless JSON before
  361. // this point, so this payload is structured-cloneable by contract.
  362. worker.postMessage(payload)
  363. }
  364. const record = bindings.get(message.global)
  365. // Own-property lookup only: a forged name like 'constructor' or
  366. // 'hasOwnProperty' must not walk the record's prototype chain and
  367. // reach a callable the consumer never declared.
  368. const fn = record && Object.hasOwn(record, message.name) ? record[message.name] : undefined
  369. if (typeof fn !== 'function') {
  370. reply({ type: 'reply', id: message.id, ok: false, message: `unknown binding ${JSON.stringify(`${message.global}.${message.name}`)}` })
  371. return
  372. }
  373. void (async () => {
  374. try {
  375. const resolved = await fn(message.args)
  376. let value: CodeJsonValue | undefined
  377. try {
  378. value = snapshotJsonValue(resolved)
  379. } catch {
  380. value = undefined
  381. }
  382. if (value === undefined) {
  383. reply({ type: 'reply', id: message.id, ok: false, message: 'binding resolution must be lossless JSON' })
  384. } else {
  385. reply({ type: 'reply', id: message.id, ok: true, value })
  386. }
  387. } catch (error: unknown) {
  388. reply({ type: 'reply', id: message.id, ok: false, message: messageOf(error) })
  389. }
  390. })()
  391. }
  392. worker.on('message', (raw: unknown) => {
  393. // Parse before touching: the peer can post ANY shape, and a throw in
  394. // this listener would crash the host process. Junk drops silently.
  395. const message = parseWorkerMessage(raw)
  396. if (!message) return
  397. if (message.type === 'log' && !settled && !output.admit(message.text, logs)) {
  398. finish(output.limit([...logs, ...strayLogs, message.text]))
  399. return
  400. }
  401. if (message.type === 'output-limit' && !settled) {
  402. finish(output.limit([...logs, ...strayLogs]))
  403. return
  404. }
  405. onCall(message)
  406. onDone(message)
  407. })
  408. worker.on('error', (error: Error) => {
  409. finish(output.failure([...logs, ...strayLogs], { kind: 'worker-exit', message: `worker error: ${error.message}` }))
  410. })
  411. worker.on('exit', (exitCode: number) => {
  412. finish(output.failure([...logs, ...strayLogs], { kind: 'worker-exit', message: `worker exited with code ${exitCode} before completing` }))
  413. })
  414. // The compute budget reads the worker's own measured busy time, so a
  415. // hot loop expires it no matter what dispatches are in flight, while a
  416. // program idling on a slow binding accrues nothing.
  417. const eluTimer = setInterval(() => {
  418. const elu = worker.performance.eventLoopUtilization()
  419. if (elu.active > this.config.computeMs) {
  420. finish(output.failure([...logs, ...strayLogs], { kind: 'timeout', message: `compute budget exhausted (${this.config.computeMs}ms busy)` }))
  421. }
  422. }, ELU_POLL_INTERVAL_MS)
  423. const wallTimer = setTimeout(() => {
  424. finish(output.failure([...logs, ...strayLogs], { kind: 'timeout', message: `wall-clock ceiling reached (${this.config.maxWallMs}ms)` }))
  425. }, this.config.maxWallMs)
  426. const onAbort = (): void => {
  427. finish(output.failure([...logs, ...strayLogs], { kind: 'abort', message: String(request.signal?.reason) }))
  428. }
  429. request.signal?.addEventListener('abort', onAbort, { once: true })
  430. const live: LiveRun = {
  431. worker,
  432. finished,
  433. settle: (failure: CodeRunFailure) => { finish(output.failure([...logs, ...strayLogs], failure)) },
  434. }
  435. this.live.add(live)
  436. })
  437. }
  438. }
  439. export default WorkerCodeRuntime