runtime.spec.ts 31 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715
  1. import { spawn } from 'node:child_process'
  2. import { mkdtemp, rm, writeFile } from 'node:fs/promises'
  3. import { tmpdir } from 'node:os'
  4. import { join } from 'node:path'
  5. import { PassThrough, Writable } from 'node:stream'
  6. import { Context } from 'cordis'
  7. import { describe, expect, it, vi } from 'vitest'
  8. import type { Sandbox } from '@deepseek-ai/dsh-e2b'
  9. import {
  10. E2BFrameDecoder,
  11. encodeE2BFrame,
  12. } from '@deepseek-ai/dsh-e2b'
  13. import type E2BSandboxService from '@deepseek-ai/dsh-e2b'
  14. import type {
  15. SubprocessHandle,
  16. SubprocessOutcome,
  17. SubprocessSpawnSpec,
  18. } from '@deepseek-ai/dsh-subprocess'
  19. import E2BSubprocessService from '@deepseek-ai/dsh-subprocess-e2b'
  20. import {
  21. encodeWorkerJson,
  22. } from '@deepseek-ai/dsh-code-runtime-worker'
  23. import E2BCodeRuntime from '@deepseek-ai/dsh-code-runtime-e2b'
  24. import * as E2BCodeRuntimeInvariant from '../src/invariant.ts'
  25. import { CODE_RUNNER_SOURCE } from '../src/runner-source.ts'
  26. import InvariantService from '@deepseek-ai/dsh-invariants'
  27. class FakeHandle implements SubprocessHandle {
  28. readonly pid = 123
  29. readonly stdin: Writable | undefined
  30. readonly stdout: PassThrough | undefined
  31. readonly stderr = undefined
  32. readonly collected: SubprocessHandle['collected']
  33. readonly done: Promise<SubprocessOutcome>
  34. readonly writes: unknown[] = []
  35. readonly result = Promise.withResolvers<SubprocessOutcome>()
  36. terminated = 0
  37. waitCalls = 0
  38. private readonly decoder = new E2BFrameDecoder(10_000_000)
  39. private readonly waitError: Error | undefined
  40. private readonly waitResult: Promise<boolean> | undefined
  41. private settled = false
  42. constructor(
  43. private readonly onMessage: (message: unknown, handle: FakeHandle) => void = () => {},
  44. options: {
  45. stdin?: boolean
  46. stdout?: boolean
  47. stderr?: string
  48. writeError?: Error
  49. waitError?: Error
  50. waitResult?: Promise<boolean>
  51. } = {},
  52. ) {
  53. this.waitError = options.waitError
  54. this.waitResult = options.waitResult
  55. this.stdin = options.stdin === false
  56. ? undefined
  57. : options.writeError === undefined
  58. ? new PassThrough()
  59. : new Writable({ write: (_chunk, _encoding, callback) => { callback(options.writeError) } })
  60. this.stdout = options.stdout === false ? undefined : new PassThrough()
  61. this.collected = options.stderr === undefined
  62. ? {}
  63. : { stderr: { readFrom: () => ({ text: options.stderr as string, nextOffset: 0, lossy: false }) } }
  64. this.done = this.result.promise
  65. this.stdin?.on('data', (chunk: Buffer) => {
  66. for (const message of this.decoder.push(chunk.toString('ascii'))) {
  67. this.writes.push(message)
  68. this.onMessage(message, this)
  69. }
  70. })
  71. }
  72. emit(message: unknown): void {
  73. this.stdout?.write(encodeE2BFrame(message))
  74. }
  75. emitRaw(text: string): void {
  76. this.stdout?.write(text)
  77. }
  78. exit(outcome: SubprocessOutcome = { exitCode: 0, signal: null }): void {
  79. if (this.settled) return
  80. this.settled = true
  81. this.stdout?.end()
  82. this.result.resolve(outcome)
  83. }
  84. crash(error: unknown): void {
  85. if (this.settled) return
  86. this.settled = true
  87. this.stdout?.end()
  88. this.result.reject(error)
  89. }
  90. terminate(): void {
  91. this.terminated += 1
  92. this.exit({ exitCode: null, signal: 'SIGTERM' })
  93. }
  94. async waitForExit(): Promise<boolean> {
  95. this.waitCalls += 1
  96. if (this.waitError !== undefined) throw this.waitError
  97. if (this.waitResult !== undefined) return await this.waitResult
  98. return true
  99. }
  100. }
  101. interface RuntimeFixture {
  102. ctx: Context
  103. fiber: Awaited<ReturnType<Context['plugin']>>
  104. runtime: E2BCodeRuntime
  105. sandbox: Sandbox
  106. spawn: ReturnType<typeof vi.fn<(spec: SubprocessSpawnSpec) => SubprocessHandle>>
  107. write: ReturnType<typeof vi.fn>
  108. run: ReturnType<typeof vi.fn>
  109. }
  110. async function setup(
  111. handles: FakeHandle[] = [],
  112. config: Record<string, number> = {},
  113. sandboxOverrides: Partial<Sandbox> = {},
  114. getSandbox?: () => Promise<Sandbox>,
  115. ): Promise<RuntimeFixture> {
  116. const write = vi.fn().mockResolvedValue([])
  117. const run = vi.fn().mockImplementation(async (command: string) => ({
  118. exitCode: 0,
  119. stdout: command.startsWith('command -v') ? '/usr/bin/node\n' : '',
  120. stderr: '',
  121. }))
  122. const sandbox = {
  123. files: { write },
  124. commands: { run },
  125. ...sandboxOverrides,
  126. } as unknown as Sandbox
  127. const e2b = {
  128. cwd: '/workspace',
  129. runtimeRoot: '/workspace/.dsh-e2b',
  130. getSandbox: getSandbox ?? (async () => sandbox),
  131. } as unknown as E2BSandboxService
  132. const spawn = vi.fn<(spec: SubprocessSpawnSpec) => SubprocessHandle>(() => {
  133. const handle = handles.shift()
  134. if (handle === undefined) throw new Error('no fake handle queued')
  135. return handle
  136. })
  137. const subprocess = Object.create(E2BSubprocessService.prototype) as E2BSubprocessService
  138. Object.defineProperty(subprocess, 'spawn', { value: spawn })
  139. const ctx = new Context()
  140. ctx.provide('e2b', e2b)
  141. ctx.provide('subprocess', subprocess)
  142. const fiber = await ctx.plugin(E2BCodeRuntime, config)
  143. return { ctx, fiber, runtime: ctx.codeRuntime as E2BCodeRuntime, sandbox, spawn, write, run }
  144. }
  145. function request(program = 'return 1') {
  146. return { program, bindings: [] }
  147. }
  148. async function runInstalledRunner(
  149. code: string,
  150. maxOutputBytes = 2_000_000,
  151. ): Promise<{ messages: unknown[]; stderr: string }> {
  152. const directory = await mkdtemp(join(tmpdir(), 'dsh-e2b-code-runner-'))
  153. const runner = join(directory, 'runner.mjs')
  154. await writeFile(runner, CODE_RUNNER_SOURCE)
  155. const child = spawn(process.execPath, [runner], { stdio: ['pipe', 'pipe', 'pipe'] })
  156. const decoder = new E2BFrameDecoder(4_000_000)
  157. const messages: unknown[] = []
  158. let stderr = ''
  159. let outputError: unknown
  160. child.stdout.setEncoding('ascii')
  161. child.stdout.on('data', (chunk: string) => {
  162. try {
  163. messages.push(...decoder.push(chunk))
  164. } catch (error: unknown) {
  165. outputError = error
  166. child.kill('SIGKILL')
  167. }
  168. })
  169. child.stderr.setEncoding('utf8')
  170. child.stderr.on('data', (chunk: string) => { stderr += chunk })
  171. try {
  172. child.stdin.write(encodeE2BFrame({
  173. type: 'boot',
  174. code,
  175. namespaces: [],
  176. computeMs: 1_000,
  177. maxOutputBytes,
  178. maxOldGenerationSizeMb: 128,
  179. }))
  180. await new Promise<void>((resolve, reject) => {
  181. const timeout = setTimeout(() => {
  182. child.kill('SIGKILL')
  183. reject(new Error('installed E2B code runner did not exit'))
  184. }, 5_000)
  185. child.once('error', (error) => {
  186. clearTimeout(timeout)
  187. reject(error)
  188. })
  189. child.once('exit', () => {
  190. clearTimeout(timeout)
  191. resolve()
  192. })
  193. })
  194. if (outputError !== undefined) throw outputError
  195. decoder.finish()
  196. return { messages, stderr }
  197. } finally {
  198. child.kill('SIGKILL')
  199. await rm(directory, { recursive: true, force: true })
  200. }
  201. }
  202. describe('E2BCodeRuntime', () => {
  203. it('keeps model-owned descriptors outside the host framing process', async () => {
  204. const forged = Buffer.from(JSON.stringify({ type: 'done' })).toString('base64') + '\\n'
  205. const { messages, stderr } = await runInstalledRunner(
  206. `
  207. const fs = await import('node:fs')
  208. const childProcess = await import('node:child_process')
  209. fs.writeSync(1, ${JSON.stringify(forged)})
  210. childProcess.spawnSync(process.execPath, ['-e', 'process.stdout.write("child-native")'], { stdio: 'inherit' })
  211. return true
  212. `,
  213. )
  214. const records = messages as Array<{ type?: string; text?: string; value?: unknown }>
  215. const terminal = records.filter(message => message.type === 'done')
  216. expect(stderr).toBe('')
  217. expect(terminal).toEqual([{ type: 'done', value: [true] }])
  218. expect(records.at(-1)).toEqual(terminal[0])
  219. expect(records.filter(message => message.type === 'log').map(message => message.text).join(''))
  220. .toContain(forged + 'child-native')
  221. })
  222. it('bounds native descriptor output before it reaches the host protocol', async () => {
  223. const { messages, stderr } = await runInstalledRunner(
  224. "(await import('node:fs')).writeSync(1, 'x'.repeat(4096)); return true",
  225. 64,
  226. )
  227. const records = messages as Array<{ type?: string; text?: string }>
  228. expect(stderr).toBe('')
  229. expect(records.at(-1)).toEqual({ type: 'output-limit' })
  230. expect(Buffer.byteLength(records.filter(message => message.type === 'log').map(message => message.text).join('')))
  231. .toBeLessThanOrEqual(62)
  232. })
  233. it('drains native worker pipes before emitting the terminal frame', async () => {
  234. const expectedBytes = 1_048_576
  235. const { messages, stderr } = await runInstalledRunner(
  236. `
  237. let stdoutPrototype = Object.getPrototypeOf(process.stdout)
  238. while (stdoutPrototype && !Object.hasOwn(stdoutPrototype, 'write')) stdoutPrototype = Object.getPrototypeOf(stdoutPrototype)
  239. Reflect.apply(stdoutPrototype.write, process.stdout, ['x'.repeat(${expectedBytes})])
  240. return true
  241. `,
  242. )
  243. const records = messages as Array<{ type?: string; text?: string }>
  244. const terminalIndex = records.findIndex(message => message.type === 'done')
  245. const nativeOutput = records
  246. .slice(0, terminalIndex)
  247. .filter(message => message.type === 'log')
  248. .map(message => message.text ?? '')
  249. .join('')
  250. expect(stderr).toBe('')
  251. expect(terminalIndex).toBe(records.length - 1)
  252. expect(Buffer.byteLength(nativeOutput)).toBe(expectedBytes)
  253. })
  254. it('prepares the remote runner and returns logs and a lossless completion', async () => {
  255. const handle = new FakeHandle((message, current) => {
  256. if ((message as { type?: string }).type !== 'boot') return
  257. current.emit({ type: 'log', text: 'remote 你好' })
  258. current.emitRaw(
  259. encodeE2BFrame({ type: 'done', value: encodeWorkerJson({ answer: 42 }) })
  260. + encodeE2BFrame({ type: 'log', text: 'ignored after done' }),
  261. )
  262. current.emit({ type: 'log', text: 'also ignored after done' })
  263. })
  264. const fixture = await setup([handle])
  265. await expect(fixture.runtime.run(request('const answer: number = 42; return { answer }')))
  266. .resolves.toEqual({ logs: ['remote 你好'], value: { answer: 42 } })
  267. expect(fixture.runtime.language).toBe('typescript')
  268. expect(fixture.runtime.isolation).toBe('container')
  269. expect(fixture.write).toHaveBeenCalledWith([{ path: '/workspace/.dsh-e2b/code-runtime-runner.mjs', data: CODE_RUNNER_SOURCE }])
  270. expect(fixture.run).toHaveBeenCalledWith("chmod 600 -- '/workspace/.dsh-e2b/code-runtime-runner.mjs'")
  271. expect(fixture.spawn).toHaveBeenCalledWith(expect.objectContaining({
  272. argv: ['/usr/bin/node', '/workspace/.dsh-e2b/code-runtime-runner.mjs'],
  273. cwd: '/workspace',
  274. env: {},
  275. }))
  276. expect(handle.terminated).toBe(1)
  277. expect(handle.waitCalls).toBe(1)
  278. await fixture.fiber.dispose()
  279. })
  280. it('bridges binding success, host rejection, unknown members, and invalid values', async () => {
  281. const replies: unknown[] = []
  282. const handle = new FakeHandle((message, current) => {
  283. const record = message as { type?: string; id?: number; ok?: boolean }
  284. if (record.type === 'boot') {
  285. current.emit({ type: 'call', id: 1, global: 'bridge', name: 'double', args: encodeWorkerJson({ value: 4 }) })
  286. current.emit({ type: 'call', id: 2, global: 'bridge', name: 'fail', args: encodeWorkerJson(null) })
  287. current.emit({ type: 'call', id: 3, global: 'bridge', name: 'missing', args: encodeWorkerJson(null) })
  288. current.emit({ type: 'call', id: 4, global: 'bridge', name: 'double', args: [] })
  289. current.emit({ type: 'call', id: 5, global: 'bridge', name: 'invalid', args: encodeWorkerJson(null) })
  290. current.emit({ type: 'call', id: 6, global: 'bridge', name: 'throwing', args: encodeWorkerJson(null) })
  291. current.emit({ type: 'call', id: 1, global: 'bridge', name: 'double', args: encodeWorkerJson({ value: 99 }) })
  292. return
  293. }
  294. if (record.type === 'reply') {
  295. replies.push(message)
  296. if (replies.length === 6) current.emit({ type: 'done', value: encodeWorkerJson('done') })
  297. }
  298. })
  299. const fixture = await setup([handle])
  300. const result = await fixture.runtime.run({
  301. program: 'return await bridge.double({ value: 4 })',
  302. bindings: [
  303. {
  304. global: 'bridge',
  305. errorClass: { name: 'BridgeError', memberNameProperty: 'member' },
  306. functions: {
  307. double: async args => (args as { value: number }).value * 2,
  308. fail: async () => { throw 'nope' },
  309. invalid: (async () => undefined) as never,
  310. throwing: async () => Object.defineProperty({}, 'value', {
  311. enumerable: true,
  312. get: () => { throw new Error('getter failed') },
  313. }),
  314. },
  315. },
  316. { global: 'plain', functions: {} },
  317. ],
  318. })
  319. expect(result).toEqual({ logs: [], value: 'done' })
  320. expect(replies.sort((left, right) => (left as { id: number }).id - (right as { id: number }).id)).toEqual([
  321. { type: 'reply', id: 1, ok: true, value: encodeWorkerJson(8) },
  322. { type: 'reply', id: 2, ok: false, message: 'nope' },
  323. { type: 'reply', id: 3, ok: false, message: 'unknown binding "bridge.missing"' },
  324. { type: 'reply', id: 4, ok: false, message: 'binding arguments must be lossless JSON' },
  325. { type: 'reply', id: 5, ok: false, message: 'binding resolution must be lossless JSON' },
  326. { type: 'reply', id: 6, ok: false, message: 'binding resolution must be lossless JSON' },
  327. ])
  328. await fixture.fiber.dispose()
  329. })
  330. it('ignores malformed runner traffic and classifies terminal runner messages', async () => {
  331. const ignored = [
  332. null, 1, {}, { type: 'log' }, { type: 'call' },
  333. { type: 'call', id: 0, global: 'x', name: 'y', args: [] },
  334. { type: 'call', id: 1, global: 1, name: 'y', args: [] },
  335. { type: 'call', id: 1, global: 'x', name: 1, args: [] },
  336. { type: 'call', id: 1, global: 'x', name: 'y', args: {} },
  337. { type: 'done', error: null },
  338. { type: 'done', error: { kind: 'invented', message: 'x' } },
  339. { type: 'done', error: { kind: 'exception', message: 1 } },
  340. ]
  341. const handles = [
  342. new FakeHandle((message, current) => {
  343. if ((message as { type?: string }).type !== 'boot') return
  344. for (const item of ignored) current.emit(item)
  345. current.emit({ type: 'done' })
  346. }),
  347. new FakeHandle((message, current) => {
  348. if ((message as { type?: string }).type === 'boot') current.emit({ type: 'done', error: { kind: 'exception', message: 'boom' } })
  349. }),
  350. new FakeHandle((message, current) => {
  351. if ((message as { type?: string }).type === 'boot') current.emit({ type: 'done', value: [] })
  352. }),
  353. new FakeHandle((message, current) => {
  354. if ((message as { type?: string }).type === 'boot') current.emit({ type: 'output-limit' })
  355. }),
  356. ]
  357. const fixture = await setup(handles, { maxOutputBytes: 64, maxFrameBytes: 128 })
  358. await expect(fixture.runtime.run(request())).resolves.toEqual({ logs: [] })
  359. await expect(fixture.runtime.run(request())).resolves.toEqual({ logs: [], error: { kind: 'exception', message: 'boom' } })
  360. await expect(fixture.runtime.run(request())).resolves.toEqual({ logs: [], error: { kind: 'invalid-output', message: 'program completion must be lossless JSON' } })
  361. await expect(fixture.runtime.run(request())).resolves.toEqual({ logs: [], error: { kind: 'output-limit', message: 'outer output exceeded 64 bytes' } })
  362. await fixture.fiber.dispose()
  363. })
  364. it('enforces the host output ledger and catches malformed bridge output', async () => {
  365. const handles = [
  366. new FakeHandle((message, current) => {
  367. if ((message as { type?: string }).type === 'boot') current.emit({ type: 'log', text: 'x'.repeat(1_000) })
  368. }),
  369. new FakeHandle((message, current) => {
  370. if ((message as { type?: string }).type === 'boot') current.emitRaw('not-base64\n')
  371. }),
  372. new FakeHandle((message, current) => {
  373. if ((message as { type?: string }).type === 'boot') current.emitRaw('é')
  374. }),
  375. new FakeHandle((message, current) => {
  376. if ((message as { type?: string }).type === 'boot') current.stdout?.emit('error', new Error('stdout broke'))
  377. }),
  378. ]
  379. const fixture = await setup(handles, { maxOutputBytes: 128, maxFrameBytes: 4_096 })
  380. expect((await fixture.runtime.run(request())).error?.kind).toBe('output-limit')
  381. const malformed = (await fixture.runtime.run(request())).error
  382. expect(malformed?.kind).toBe('worker-exit')
  383. expect(malformed?.message).toContain('bridge failed')
  384. expect((await fixture.runtime.run(request())).error?.message).toContain('non-ASCII')
  385. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime stdout failed: stdout broke' })
  386. await fixture.fiber.dispose()
  387. })
  388. it('enforces the outbound frame bound on boot and binding replies', async () => {
  389. const oversizedBoot = new FakeHandle()
  390. const oversizedReply = new FakeHandle((message, current) => {
  391. if ((message as { type?: string }).type === 'boot') {
  392. current.emit({ type: 'call', id: 1, global: 'bridge', name: 'large', args: encodeWorkerJson(null) })
  393. }
  394. })
  395. const fixture = await setup([oversizedBoot, oversizedReply], { maxOutputBytes: 128, maxFrameBytes: 512 })
  396. const bootResult = await fixture.runtime.run(request(`return ${JSON.stringify('x'.repeat(1_000))}`))
  397. expect(bootResult.error).toMatchObject({ kind: 'worker-exit' })
  398. expect(bootResult.error?.message).toContain('frame exceeded its byte limit')
  399. expect(oversizedBoot.writes).toHaveLength(0)
  400. const replyResult = await fixture.runtime.run({
  401. program: 'return await bridge.large(null)',
  402. bindings: [{ global: 'bridge', functions: { large: async () => 'x'.repeat(1_000) } }],
  403. })
  404. expect(replyResult.error).toMatchObject({ kind: 'worker-exit' })
  405. expect(replyResult.error?.message).toContain('frame exceeded its byte limit')
  406. expect(oversizedReply.writes).toHaveLength(1)
  407. await fixture.fiber.dispose()
  408. })
  409. it('contains stdin errors, process exits, spawn failures, and missing pipes', async () => {
  410. const writeError = new FakeHandle(() => {}, { writeError: new Error('write callback broke') })
  411. const stdinError = new FakeHandle((message, current) => {
  412. if ((message as { type?: string }).type === 'boot') current.stdin?.emit('error', new Error('stdin broke'))
  413. })
  414. const earlyExit = new FakeHandle(() => {}, { stderr: 'remote diagnostic' })
  415. const quietExit = new FakeHandle()
  416. const emptyStderrExit = new FakeHandle(() => {}, { stderr: '' })
  417. const spawnFailure = new FakeHandle()
  418. const missingStdin = new FakeHandle(() => {}, { stdin: false })
  419. const missingStdout = new FakeHandle(() => {}, { stdout: false, waitError: new Error('missing-stream process query failed') })
  420. const truncated = new FakeHandle((message, current) => {
  421. if ((message as { type?: string }).type === 'boot') {
  422. current.emitRaw('YQ==')
  423. setImmediate(() => { current.exit() })
  424. }
  425. })
  426. const cleanupFailure = new FakeHandle((message, current) => {
  427. if ((message as { type?: string }).type === 'boot') current.emit({ type: 'done' })
  428. }, { waitError: new Error('process query failed') })
  429. const fixture = await setup([
  430. writeError, stdinError, earlyExit, quietExit, emptyStderrExit,
  431. spawnFailure, missingStdin, missingStdout, truncated, cleanupFailure,
  432. ])
  433. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime bridge write failed: write callback broke' })
  434. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime stdin failed: stdin broke' })
  435. setImmediate(() => { earlyExit.exit() })
  436. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime exited before completing: remote diagnostic' })
  437. setImmediate(() => { quietExit.exit() })
  438. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime exited before completing' })
  439. setImmediate(() => { emptyStderrExit.exit() })
  440. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime exited before completing' })
  441. setImmediate(() => { spawnFailure.crash('spawn rejected') })
  442. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime spawn failed: spawn rejected' })
  443. expect((await fixture.runtime.run(request())).error?.message).toContain('dropped a piped runtime stream')
  444. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime cleanup failed: missing-stream process query failed' })
  445. expect(missingStdin.terminated).toBe(1)
  446. expect(missingStdin.waitCalls).toBe(1)
  447. expect(missingStdout.terminated).toBe(1)
  448. expect(missingStdout.waitCalls).toBe(1)
  449. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B frame stream ended mid-frame' })
  450. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'worker-exit', message: 'E2B runtime cleanup failed: process query failed' })
  451. await fixture.fiber.dispose()
  452. })
  453. it('reports wall timeout, abort, pre-abort, type-strip failure, and disposal', async () => {
  454. const timeout = new FakeHandle()
  455. const abort = new FakeHandle()
  456. const disposing = new FakeHandle()
  457. const fixture = await setup([timeout, abort, disposing], { maxWallMs: 20 })
  458. expect((await fixture.runtime.run(request())).error).toEqual({ kind: 'timeout', message: 'wall-clock ceiling reached (20ms)' })
  459. const controller = new AbortController()
  460. const aborting = fixture.runtime.run({ ...request(), signal: controller.signal })
  461. controller.abort('stop')
  462. expect((await aborting).error).toEqual({ kind: 'abort', message: 'stop' })
  463. expect((await fixture.runtime.run({ ...request(), signal: AbortSignal.abort('already') })).error)
  464. .toEqual({ kind: 'abort', message: 'already' })
  465. expect((await fixture.runtime.run(request('enum E { A }'))).error?.kind).toBe('exception')
  466. const live = fixture.runtime.run(request())
  467. await new Promise(resolve => setImmediate(resolve))
  468. await fixture.fiber.dispose()
  469. expect((await live).error).toEqual({ kind: 'abort', message: 'runtime disposed' })
  470. await expect(fixture.runtime.run(request())).rejects.toThrow('after disposal')
  471. })
  472. it('drops binding replies that settle after abort', async () => {
  473. const controller = new AbortController()
  474. const resolution = Promise.withResolvers<string>()
  475. const invoked = Promise.withResolvers<undefined>()
  476. const handle = new FakeHandle((message, current) => {
  477. if ((message as { type?: string }).type === 'boot') {
  478. current.emit({ type: 'call', id: 1, global: 'bridge', name: 'late', args: encodeWorkerJson(null) })
  479. }
  480. })
  481. const fixture = await setup([handle])
  482. const running = fixture.runtime.run({
  483. program: 'return await bridge.late(null)',
  484. bindings: [{
  485. global: 'bridge',
  486. functions: {
  487. late: async () => {
  488. invoked.resolve(undefined)
  489. return await resolution.promise
  490. },
  491. },
  492. }],
  493. signal: controller.signal,
  494. })
  495. await invoked.promise
  496. controller.abort('stop')
  497. expect((await running).error).toEqual({ kind: 'abort', message: 'stop' })
  498. resolution.resolve('late')
  499. await new Promise(resolve => setImmediate(resolve))
  500. expect(handle.writes).toHaveLength(1)
  501. await fixture.fiber.dispose()
  502. })
  503. it('validates binding and runtime configuration before remote execution', async () => {
  504. const fixture = await setup([])
  505. const invalidRequests = [
  506. { global: 'not-valid!', functions: {} },
  507. { global: 'await', functions: {} },
  508. { global: 'console', functions: {} },
  509. { global: 'same', functions: {} },
  510. { global: 'same', functions: {} },
  511. { global: 'ok', functions: {}, errorClass: { name: 'not-valid!', memberNameProperty: 'member' } },
  512. { global: 'ok', functions: {}, errorClass: { name: 'await', memberNameProperty: 'member' } },
  513. { global: 'Clash', functions: {}, errorClass: { name: 'Clash', memberNameProperty: 'member' } },
  514. { global: 'one', functions: {}, errorClass: { name: 'Err', memberNameProperty: 'member' } },
  515. { global: 'two', functions: {}, errorClass: { name: 'Err', memberNameProperty: 'member' } },
  516. { global: 'ok', functions: {}, errorClass: { name: 'Err', memberNameProperty: '' } },
  517. { global: 'ok', functions: {}, errorClass: { name: 'Err', memberNameProperty: 'message' } },
  518. ]
  519. for (const bindings of [
  520. [invalidRequests[0]], [invalidRequests[1]], [invalidRequests[2]],
  521. invalidRequests.slice(3, 5), [invalidRequests[5]], [invalidRequests[6]],
  522. [invalidRequests[7]], invalidRequests.slice(8, 10), [invalidRequests[10]], [invalidRequests[11]],
  523. ]) {
  524. await expect(fixture.runtime.run({ program: 'return 1', bindings: bindings as never })).rejects.toThrow()
  525. }
  526. await fixture.fiber.dispose()
  527. for (const config of [
  528. { computeMs: 0 }, { computeMs: 1.5 }, { maxOutputBytes: 3 },
  529. { maxWallMs: 2_147_483_648 }, { maxFrameBytes: 10, maxOutputBytes: 20 },
  530. ]) {
  531. const ctx = new Context()
  532. const subprocess = Object.create(E2BSubprocessService.prototype) as E2BSubprocessService
  533. ctx.provide('e2b', { getSandbox: async () => ({}) } as never)
  534. ctx.provide('subprocess', subprocess)
  535. await expect(ctx.plugin(E2BCodeRuntime, config)).rejects.toThrow()
  536. }
  537. const wrong = new Context()
  538. wrong.provide('e2b', { getSandbox: async () => ({}) } as never)
  539. wrong.provide('subprocess', {} as never)
  540. await expect(wrong.plugin(E2BCodeRuntime, {})).rejects.toThrow('dsh-subprocess-e2b')
  541. })
  542. it('turns asynchronous runtime preparation failure into a run result', async () => {
  543. const sandbox = {
  544. files: { write: vi.fn().mockRejectedValue(new Error('upload failed')) },
  545. commands: { run: vi.fn() },
  546. } as unknown as Sandbox
  547. const fixture = await setup([], {}, sandbox)
  548. expect((await fixture.runtime.run(request())).error).toEqual({
  549. kind: 'worker-exit',
  550. message: 'E2B runtime setup failed: upload failed',
  551. })
  552. await fixture.fiber.dispose()
  553. })
  554. it('returns disposal when remote preparation completes after teardown', async () => {
  555. const gate = Promise.withResolvers<Sandbox>()
  556. const fixture = await setup([], {}, {}, () => gate.promise)
  557. const running = fixture.runtime.run(request())
  558. const disposing = fixture.fiber.dispose()
  559. let disposed = false
  560. void disposing.then(() => { disposed = true })
  561. await new Promise(resolve => setImmediate(resolve))
  562. const disposedBeforeSetup = disposed
  563. gate.resolve(fixture.sandbox)
  564. await disposing
  565. expect(disposedBeforeSetup).toBe(false)
  566. expect((await running).error).toEqual({ kind: 'abort', message: 'runtime disposed' })
  567. expect(fixture.write).not.toHaveBeenCalled()
  568. })
  569. it('observes abort while runtime preparation is pending', async () => {
  570. const gate = Promise.withResolvers<Sandbox>()
  571. const fixture = await setup([], {}, {}, () => gate.promise)
  572. const controller = new AbortController()
  573. const running = fixture.runtime.run({ ...request(), signal: controller.signal })
  574. controller.abort('stop during setup')
  575. const early = await Promise.race([
  576. running.then(result => ({ kind: 'result' as const, result })),
  577. new Promise<{ kind: 'pending' }>((resolve) => { setImmediate(() => { resolve({ kind: 'pending' }) }) }),
  578. ])
  579. expect(fixture.spawn).not.toHaveBeenCalled()
  580. gate.resolve(fixture.sandbox)
  581. expect(early).toMatchObject({ kind: 'result', result: { error: { kind: 'abort', message: 'stop during setup' } } })
  582. await running
  583. await fixture.fiber.dispose()
  584. })
  585. it('classifies an abort that races synchronous subprocess spawn', async () => {
  586. const fixture = await setup()
  587. const controller = new AbortController()
  588. fixture.spawn.mockImplementationOnce(() => {
  589. controller.abort('stop at spawn')
  590. throw new Error('aborted before spawn')
  591. })
  592. expect((await fixture.runtime.run({ ...request(), signal: controller.signal })).error)
  593. .toEqual({ kind: 'abort', message: 'stop at spawn' })
  594. fixture.spawn.mockImplementationOnce(() => { throw new Error('synchronous spawn failure') })
  595. expect((await fixture.runtime.run(request())).error).toEqual({
  596. kind: 'worker-exit',
  597. message: 'E2B runtime spawn failed: synchronous spawn failure',
  598. })
  599. await fixture.fiber.dispose()
  600. const disposingFixture = await setup()
  601. disposingFixture.spawn.mockImplementationOnce(() => {
  602. void (disposingFixture.runtime as unknown as { teardown(): Promise<void> }).teardown()
  603. throw new Error('spawn raced disposal')
  604. })
  605. expect((await disposingFixture.runtime.run(request())).error)
  606. .toEqual({ kind: 'abort', message: 'runtime disposed' })
  607. await disposingFixture.fiber.dispose()
  608. })
  609. it('closes both abort races around runtime readiness and live-run publication', async () => {
  610. let preparationAborted = false
  611. const preparationSignal = {
  612. get aborted() { return preparationAborted },
  613. reason: 'preparation race',
  614. addEventListener() { preparationAborted = true },
  615. removeEventListener() {},
  616. } as unknown as AbortSignal
  617. const liveHandle = new FakeHandle()
  618. const fixture = await setup([liveHandle])
  619. expect((await fixture.runtime.run({ ...request(), signal: preparationSignal })).error)
  620. .toEqual({ kind: 'abort', message: 'preparation race' })
  621. expect(fixture.spawn).not.toHaveBeenCalled()
  622. let liveAborted = false
  623. let registrations = 0
  624. const liveSignal = {
  625. get aborted() { return liveAborted },
  626. reason: 'live publication race',
  627. addEventListener() {
  628. registrations += 1
  629. if (registrations === 2) liveAborted = true
  630. },
  631. removeEventListener() {},
  632. } as unknown as AbortSignal
  633. expect((await fixture.runtime.run({ ...request(), signal: liveSignal })).error)
  634. .toEqual({ kind: 'abort', message: 'live publication race' })
  635. await fixture.fiber.dispose()
  636. })
  637. it('retains a live run until remote cleanup reaches quiescence', async () => {
  638. const cleanup = Promise.withResolvers<boolean>()
  639. const handle = new FakeHandle((message, current) => {
  640. if ((message as { type?: string }).type === 'boot') current.emit({ type: 'done' })
  641. }, { waitResult: cleanup.promise })
  642. const fixture = await setup([handle])
  643. const running = fixture.runtime.run(request())
  644. await vi.waitFor(() => { expect(handle.waitCalls).toBe(1) })
  645. const disposing = fixture.fiber.dispose()
  646. let disposed = false
  647. void disposing.then(() => { disposed = true })
  648. await new Promise(resolve => setImmediate(resolve))
  649. const disposedBeforeCleanup = disposed
  650. cleanup.resolve(true)
  651. await expect(running).resolves.toEqual({ logs: [] })
  652. await expect(disposing).resolves.toBeUndefined()
  653. expect(disposedBeforeCleanup).toBe(false)
  654. })
  655. it('registers the package-owned invariant companion', async () => {
  656. const ctx = new Context()
  657. await ctx.plugin(InvariantService, { enabled: true })
  658. const fiber = await ctx.plugin(E2BCodeRuntimeInvariant).await()
  659. await fiber.dispose()
  660. })
  661. })