1
0

runtime-host.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251
  1. /** Shared host mechanics for local and subprocess-hosted TypeScript worker runtimes. */
  2. import { stripTypeScriptTypes } from 'node:module'
  3. import type { Readable } from 'node:stream'
  4. import type {
  5. CodeBindingNamespace,
  6. CodeJsonValue,
  7. CodeRunFailure,
  8. CodeRunRequest,
  9. CodeRunResult,
  10. } from '@deepseek-ai/dsh-code-runtime'
  11. import { jsonStringBytesUpTo, jsonValueBytesUpTo, truncateJsonStringBytes } from './output-json.ts'
  12. import { decodeWorkerJson, encodeWorkerJson, snapshotCodeJsonValue } from './worker-json.ts'
  13. import type { WorkerJsonWire } from './worker-json.ts'
  14. /** Smallest cap that can represent an empty log array and failure message. */
  15. export const MIN_RUNTIME_OUTPUT_BYTES = 4
  16. /**
  17. * Resolve after a worker pipe emits queued data or closes during termination.
  18. * @param stream - captured worker or child-process pipe.
  19. * @returns after no more queued bytes can arrive.
  20. */
  21. export function waitForRuntimePipeDrain(stream: Readable): Promise<void> {
  22. if (stream.readableEnded || stream.destroyed) return Promise.resolve()
  23. return new Promise((resolve) => {
  24. const done = (): void => {
  25. stream.off('end', done)
  26. stream.off('close', done)
  27. stream.off('error', done)
  28. resolve()
  29. }
  30. stream.once('end', done)
  31. stream.once('close', done)
  32. stream.once('error', done)
  33. /* v8 ignore next -- termination can win the adjacent listener-registration race. */
  34. if (stream.readableEnded || stream.destroyed) done()
  35. })
  36. }
  37. const IDENTIFIER = /^[A-Za-z_$][A-Za-z0-9_$]*$/
  38. const RESERVED_WORDS = new Set([
  39. 'await', 'break', 'case', 'catch', 'class', 'const', 'continue', 'debugger', 'default', 'delete', 'do',
  40. 'else', 'enum', 'export', 'extends', 'false', 'finally', 'for', 'function', 'if', 'import', 'in',
  41. 'instanceof', 'new', 'null', 'return', 'super', 'switch', 'this', 'throw', 'true', 'try', 'typeof',
  42. 'var', 'void', 'while', 'with', 'yield', 'let', 'static', 'implements', 'interface', 'package',
  43. 'private', 'protected', 'public', 'arguments', 'eval',
  44. ])
  45. const RESERVED_ERROR_PROPERTIES = new Set(['name', 'message', 'stack'])
  46. const STRIP_WRAP = { prefix: 'async function __dsh_program__() {\n', suffix: '\n}' } as const
  47. /** One validated binding call received from an isolated worker. */
  48. export interface RuntimeBindingCall {
  49. /** Correlation id supplied by the isolated worker. */
  50. readonly id: number
  51. /** Injected namespace global. */
  52. readonly global: string
  53. /** Declared namespace function. */
  54. readonly name: string
  55. /** Untrusted lossless-JSON wire payload. */
  56. readonly args: unknown
  57. }
  58. /** One host reply to an isolated worker binding call. */
  59. export type RuntimeBindingReply =
  60. | { readonly type: 'reply'; readonly id: number; readonly ok: true; readonly value: WorkerJsonWire }
  61. | { readonly type: 'reply'; readonly id: number; readonly ok: false; readonly message: string }
  62. /**
  63. * Render an unknown thrown value without assuming it is an Error.
  64. * @param error - thrown or rejected value.
  65. * @returns the caller-facing diagnostic text.
  66. */
  67. export function runtimeErrorMessage(error: unknown): string {
  68. try {
  69. return error instanceof Error ? error.message : String(error)
  70. } catch {
  71. return 'binding rejected with an unrenderable value'
  72. }
  73. }
  74. /**
  75. * Strip erasable TypeScript while preserving the program's body coordinates.
  76. * @param program - model-written async-function body.
  77. * @returns JavaScript source with the wrapper removed.
  78. */
  79. export function stripRuntimeProgram(program: string): string {
  80. const stripped = stripTypeScriptTypes(STRIP_WRAP.prefix + program + STRIP_WRAP.suffix)
  81. return stripped.slice(STRIP_WRAP.prefix.length, stripped.length - STRIP_WRAP.suffix.length)
  82. }
  83. /**
  84. * Validate binding globals and typed-error declarations shared by worker runtimes.
  85. * @param request - code-runtime request carrying the namespaces.
  86. * @param implementationName - package name used in seam-misuse diagnostics.
  87. * @returns namespaces indexed by their injected global.
  88. */
  89. export function validateRuntimeBindings(
  90. request: CodeRunRequest,
  91. implementationName: string,
  92. ): Map<string, CodeBindingNamespace> {
  93. const bindings = new Map<string, CodeBindingNamespace>()
  94. for (const namespace of request.bindings) {
  95. if (!IDENTIFIER.test(namespace.global) || RESERVED_WORDS.has(namespace.global)) {
  96. throw new Error(`${implementationName}: binding global ${JSON.stringify(namespace.global)} is not a usable identifier`)
  97. }
  98. if (namespace.global === 'console' || bindings.has(namespace.global)) {
  99. throw new Error(`${implementationName}: duplicate binding global ${JSON.stringify(namespace.global)}`)
  100. }
  101. bindings.set(namespace.global, namespace)
  102. }
  103. const errorClassNames = new Set<string>()
  104. for (const namespace of request.bindings) {
  105. const descriptor = namespace.errorClass
  106. if (descriptor === undefined) continue
  107. if (!IDENTIFIER.test(descriptor.name) || RESERVED_WORDS.has(descriptor.name)) {
  108. throw new Error(`${implementationName}: binding error class ${JSON.stringify(descriptor.name)} is not a usable identifier`)
  109. }
  110. if (descriptor.name === 'console' || bindings.has(descriptor.name) || errorClassNames.has(descriptor.name)) {
  111. throw new Error(`${implementationName}: duplicate injected global ${JSON.stringify(descriptor.name)}`)
  112. }
  113. if (descriptor.memberNameProperty.length === 0 || RESERVED_ERROR_PROPERTIES.has(descriptor.memberNameProperty)) {
  114. throw new Error(`${implementationName}: binding error member property ${JSON.stringify(descriptor.memberNameProperty)} is not usable`)
  115. }
  116. errorClassNames.add(descriptor.name)
  117. }
  118. return bindings
  119. }
  120. /**
  121. * Resolve one untrusted worker call through a declared host binding.
  122. * @param call - parsed call envelope from the isolated worker.
  123. * @param bindings - namespaces returned by {@link validateRuntimeBindings}.
  124. * @returns a lossless-JSON success or stable rejection reply.
  125. */
  126. export async function invokeRuntimeBinding(
  127. call: RuntimeBindingCall,
  128. bindings: ReadonlyMap<string, CodeBindingNamespace>,
  129. ): Promise<RuntimeBindingReply> {
  130. const functions = bindings.get(call.global)?.functions
  131. const fn = functions !== undefined && Object.hasOwn(functions, call.name) ? functions[call.name] : undefined
  132. if (typeof fn !== 'function') {
  133. return { type: 'reply', id: call.id, ok: false, message: `unknown binding ${JSON.stringify(`${call.global}.${call.name}`)}` }
  134. }
  135. const args = decodeWorkerJson(call.args)
  136. if (args === undefined) {
  137. return { type: 'reply', id: call.id, ok: false, message: 'binding arguments must be lossless JSON' }
  138. }
  139. try {
  140. const resolved = await fn(args)
  141. let value: CodeJsonValue | undefined
  142. try {
  143. value = snapshotCodeJsonValue(resolved)
  144. } catch {
  145. value = undefined
  146. }
  147. if (value === undefined) {
  148. return { type: 'reply', id: call.id, ok: false, message: 'binding resolution must be lossless JSON' }
  149. }
  150. return { type: 'reply', id: call.id, ok: true, value: encodeWorkerJson(value) }
  151. } catch (error: unknown) {
  152. return { type: 'reply', id: call.id, ok: false, message: runtimeErrorMessage(error) }
  153. }
  154. }
  155. /** One run's combined outer-output ledger; binding values never enter it. */
  156. export class RuntimeOutputLedger {
  157. private bytes = 2
  158. private entries = 0
  159. /** @param maxBytes - hard cap for logs plus completion or failure payload. */
  160. constructor(private readonly maxBytes: number) {}
  161. /**
  162. * Admit one exact log entry.
  163. * @param text - candidate log entry.
  164. * @param sink - ordered retained log list.
  165. * @returns false when the hard cap was crossed.
  166. */
  167. admit(text: string, sink: string[]): boolean {
  168. const separatorBytes = this.entries > 0 ? 1 : 0
  169. const stringBytes = jsonStringBytesUpTo(text, this.maxBytes - this.bytes - separatorBytes)
  170. if (stringBytes === undefined) return false
  171. this.bytes += stringBytes + separatorBytes
  172. this.entries += 1
  173. sink.push(text)
  174. return true
  175. }
  176. /**
  177. * Finalize a successful completion against the combined cap.
  178. * @param logs - retained ordered logs.
  179. * @param value - optional lossless-JSON completion.
  180. * @returns the completion or output-limit result.
  181. */
  182. success(logs: string[], value?: CodeJsonValue): CodeRunResult {
  183. if (value !== undefined && jsonValueBytesUpTo(value, this.maxBytes - this.bytes) === undefined) return this.limit(logs)
  184. return { logs, ...value !== undefined ? { value } : {} }
  185. }
  186. /**
  187. * Finalize one failure diagnostic against the combined cap.
  188. * @param logs - retained ordered logs.
  189. * @param error - structured runtime failure.
  190. * @returns the failure or output-limit result.
  191. */
  192. failure(logs: string[], error: CodeRunFailure): CodeRunResult {
  193. if (jsonStringBytesUpTo(error.message, this.maxBytes - this.bytes) === undefined) return this.limit(logs)
  194. return { logs, error }
  195. }
  196. /**
  197. * Build an explicit output-limit failure with a fitting log prefix.
  198. * @param logs - ordered logs observed before the limit.
  199. * @returns bounded output-limit result.
  200. */
  201. limit(logs: string[]): CodeRunResult {
  202. const fullMessage = `outer output exceeded ${this.maxBytes} bytes`
  203. const messageBytes = fullMessage.length + 2
  204. const retained: string[] = []
  205. let retainedBytes = 2
  206. const logBudget = this.maxBytes - messageBytes
  207. for (const text of logs) {
  208. const separatorBytes = retained.length > 0 ? 1 : 0
  209. const availableBytes = logBudget - retainedBytes - separatorBytes
  210. const stringBytes = jsonStringBytesUpTo(text, availableBytes)
  211. if (stringBytes !== undefined) {
  212. retained.push(text)
  213. retainedBytes += stringBytes + separatorBytes
  214. continue
  215. }
  216. const prefix = truncateJsonStringBytes(text, availableBytes)
  217. if (prefix.length > 0) {
  218. const prefixBytes = jsonStringBytesUpTo(prefix, availableBytes)
  219. /* v8 ignore next -- truncateJsonStringBytes guarantees the same bound. */
  220. if (prefixBytes === undefined) throw new Error('output ledger produced an oversized log prefix')
  221. retained.push(prefix)
  222. retainedBytes += prefixBytes + separatorBytes
  223. }
  224. break
  225. }
  226. const message = truncateJsonStringBytes(fullMessage, this.maxBytes - retainedBytes)
  227. return { logs: retained, error: { kind: 'output-limit', message } }
  228. }
  229. }
  230. export { decodeWorkerJson, encodeWorkerJson, snapshotCodeJsonValue } from './worker-json.ts'
  231. export { jsonStringBytesUpTo, jsonValueBytesUpTo } from './output-json.ts'
  232. export { runWorkerMain } from './bootstrap.ts'
  233. export type { WorkerJsonWire } from './worker-json.ts'