index.ts 26 KB

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