subprocess.spec.ts 50 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227
  1. import { once } from 'node:events'
  2. import { Context } from 'cordis'
  3. import {
  4. CommandExitError,
  5. FileNotFoundError,
  6. type CommandHandle,
  7. type CommandResult,
  8. type Sandbox,
  9. } from '@deepseek-ai/dsh-e2b'
  10. import type E2BSandboxService from '@deepseek-ai/dsh-e2b'
  11. import type { SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'
  12. import E2BSubprocessService from '@deepseek-ai/dsh-subprocess-e2b'
  13. import * as E2BSubprocessInvariant from '../src/invariant.ts'
  14. import { E2BBase64Decoder, E2B_OUTPUT_COMPLETE_FRAME, E2BOutputReader } from '../src/output.ts'
  15. import { E2BSubprocessHandle } from '../src/process.ts'
  16. import InvariantService from '@deepseek-ai/dsh-invariants'
  17. import { describe, expect, it, vi } from 'vitest'
  18. function commandError(exitCode: number): CommandExitError {
  19. return new CommandExitError({ exitCode, stdout: '', stderr: '', error: `exit ${exitCode}` })
  20. }
  21. interface StartOptions {
  22. background: true
  23. cwd: string
  24. stdin: boolean
  25. timeoutMs: number
  26. signal?: AbortSignal
  27. envs?: Record<string, string>
  28. onStdout?: (data: string) => void | Promise<void>
  29. onStderr?: (data: string) => void | Promise<void>
  30. }
  31. class FakeCommandHandle {
  32. pid = 4242
  33. readonly sent: Array<string | Uint8Array> = []
  34. closes = 0
  35. kills = 0
  36. disconnects = 0
  37. killError: unknown
  38. disconnectError: unknown
  39. private readonly result = Promise.withResolvers<CommandResult>()
  40. private settled = false
  41. constructor(private readonly onKill: () => void = () => {}) {}
  42. wait(): Promise<CommandResult> {
  43. return this.result.promise
  44. }
  45. async sendStdin(data: string | Uint8Array): Promise<void> {
  46. this.sent.push(data)
  47. }
  48. async closeStdin(): Promise<void> {
  49. this.closes += 1
  50. }
  51. async kill(): Promise<boolean> {
  52. this.kills += 1
  53. if (this.killError !== undefined) throw this.killError
  54. this.onKill()
  55. return true
  56. }
  57. async disconnect(): Promise<void> {
  58. this.disconnects += 1
  59. if (this.disconnectError !== undefined) throw this.disconnectError
  60. }
  61. succeed(exitCode = 0): void {
  62. if (this.settled) return
  63. this.settled = true
  64. this.result.resolve({ exitCode, stdout: '', stderr: '' })
  65. }
  66. fail(exitCode: number): void {
  67. if (this.settled) return
  68. this.settled = true
  69. this.result.reject(commandError(exitCode))
  70. }
  71. crash(error: unknown): void {
  72. if (this.settled) return
  73. this.settled = true
  74. this.result.reject(error)
  75. }
  76. }
  77. class FakeSandbox {
  78. readonly handle: FakeCommandHandle
  79. readonly commandsSeen: string[] = []
  80. readonly writtenFiles: string[][] = []
  81. readonly writtenFileData = new Map<string, string>()
  82. readonly removed: string[] = []
  83. readonly directories: string[] = []
  84. startOptions: StartOptions | undefined
  85. backgroundError: unknown
  86. envError: unknown
  87. nextRemoveError: unknown
  88. probeError: unknown
  89. signalError: unknown
  90. readonly signalErrors: unknown[] = []
  91. trapsTerm = false
  92. delaysKill = false
  93. delaysKillCompletion = false
  94. sdkKillStops = true
  95. alive = true
  96. ambient = 'PATH=/ambient/bin\0KEEP=safe\0NPM_TOKEN=secret\0DSH_STALE=old\0BROKEN\0=bad\0'
  97. processGroupId = '4242\n'
  98. exitStatus = ''
  99. readonly processGroupReads: string[] = []
  100. afterStatusRead: (() => void) | undefined
  101. beforeProbe: (() => void) | undefined
  102. afterProbe: (() => void) | undefined
  103. private startGate: Promise<void> | undefined
  104. private openStart: (() => void) | undefined
  105. private processGroupReadGate: Promise<void> | undefined
  106. private openProcessGroupRead: (() => void) | undefined
  107. constructor() {
  108. this.handle = new FakeCommandHandle(() => {
  109. if (this.sdkKillStops) {
  110. this.alive = false
  111. this.handle.fail(137)
  112. }
  113. })
  114. }
  115. deferStart(): void {
  116. const gate = Promise.withResolvers<undefined>()
  117. this.startGate = gate.promise
  118. this.openStart = () => { gate.resolve(undefined) }
  119. }
  120. releaseStart(): void {
  121. this.openStart?.()
  122. }
  123. deferProcessGroupRead(): void {
  124. const gate = Promise.withResolvers<undefined>()
  125. this.processGroupReadGate = gate.promise
  126. this.openProcessGroupRead = () => { gate.resolve(undefined) }
  127. }
  128. releaseProcessGroupRead(): void {
  129. this.openProcessGroupRead?.()
  130. }
  131. finish(exitCode = 0): void {
  132. this.alive = false
  133. void this.completeOutput().then(
  134. () => {
  135. if (exitCode === 0) this.handle.succeed(0)
  136. else this.handle.fail(exitCode)
  137. },
  138. (error: unknown) => { this.handle.crash(error) },
  139. )
  140. }
  141. async completeOutput(): Promise<void> {
  142. await Promise.all([
  143. this.stdoutWire(`${E2B_OUTPUT_COMPLETE_FRAME}\n`),
  144. this.stderrWire(`${E2B_OUTPUT_COMPLETE_FRAME}\n`),
  145. ])
  146. }
  147. async stdout(data: string): Promise<void> {
  148. await this.stdoutWire(data.length === 0 ? '' : `${Buffer.from(data).toString('base64')}\n`)
  149. }
  150. async stderr(data: string): Promise<void> {
  151. await this.stderrWire(data.length === 0 ? '' : `${Buffer.from(data).toString('base64')}\n`)
  152. }
  153. async stdoutWire(data: string): Promise<void> {
  154. await this.startOptions?.onStdout?.(data)
  155. }
  156. async stderrWire(data: string): Promise<void> {
  157. await this.startOptions?.onStderr?.(data)
  158. }
  159. readonly sandbox = {
  160. sandboxId: 'fake',
  161. files: {
  162. makeDir: async (path: string): Promise<boolean> => {
  163. this.directories.push(path)
  164. return true
  165. },
  166. write: async (files: Array<{ path: string; data: string }>): Promise<object[]> => {
  167. this.writtenFiles.push(files.map(file => file.path))
  168. for (const file of files) this.writtenFileData.set(file.path, file.data)
  169. return files.map(() => ({}))
  170. },
  171. read: async (path: string): Promise<string> => {
  172. if (!path.endsWith('/exit-code')) {
  173. await this.processGroupReadGate
  174. return this.processGroupReads.shift() ?? this.processGroupId
  175. }
  176. this.afterStatusRead?.()
  177. return this.exitStatus
  178. },
  179. remove: async (path: string): Promise<void> => {
  180. this.removed.push(path)
  181. if (this.nextRemoveError !== undefined) {
  182. const error = this.nextRemoveError
  183. this.nextRemoveError = undefined
  184. throw error
  185. }
  186. },
  187. },
  188. commands: {
  189. run: async (command: string, options?: StartOptions | { signal?: AbortSignal }): Promise<CommandHandle | CommandResult> => {
  190. this.commandsSeen.push(command)
  191. if (command === 'env -0') {
  192. if (this.envError !== undefined) throw this.envError
  193. return { exitCode: 0, stdout: this.ambient, stderr: '' }
  194. }
  195. if (command.startsWith('kill -0 ')) {
  196. this.beforeProbe?.()
  197. if (options?.signal?.aborted === true) throw new DOMException('aborted', 'AbortError')
  198. if (this.probeError !== undefined) {
  199. const error = this.probeError
  200. this.probeError = undefined
  201. throw error
  202. }
  203. if (!this.alive) throw commandError(1)
  204. this.afterProbe?.()
  205. return { exitCode: 0, stdout: '', stderr: '' }
  206. }
  207. if (command.startsWith('kill -TERM ')) {
  208. const error = this.signalErrors.shift() ?? this.signalError
  209. if (error !== undefined) {
  210. if (this.signalErrors.length === 0) this.signalError = undefined
  211. throw error
  212. }
  213. if (!this.trapsTerm) {
  214. this.alive = false
  215. this.handle.fail(143)
  216. }
  217. return { exitCode: 0, stdout: '', stderr: '' }
  218. }
  219. if (command.startsWith('kill -KILL ')) {
  220. const error = this.signalErrors.shift() ?? this.signalError
  221. if (error !== undefined) {
  222. if (this.signalErrors.length === 0) this.signalError = undefined
  223. throw error
  224. }
  225. if (!this.delaysKill) this.alive = false
  226. if (!this.delaysKillCompletion) this.handle.fail(137)
  227. return { exitCode: 0, stdout: '', stderr: '' }
  228. }
  229. if ((options as StartOptions | undefined)?.background === true) {
  230. this.startOptions = options as StartOptions
  231. await this.startGate
  232. if (this.backgroundError !== undefined) throw this.backgroundError
  233. return this.handle as unknown as CommandHandle
  234. }
  235. return { exitCode: 0, stdout: '', stderr: '' }
  236. },
  237. },
  238. } as unknown as Sandbox
  239. }
  240. function spec(overrides: Partial<SubprocessSpawnSpec> = {}): SubprocessSpawnSpec {
  241. return {
  242. argv: ['bash', '-c', 'printf ok'],
  243. cwd: '/workspace',
  244. stdio: {
  245. stdin: 'ignore',
  246. stdout: { maxBytes: 4, spill: { maxBytes: 16 } },
  247. stderr: { maxBytes: 4 },
  248. },
  249. graceMs: 5,
  250. ...overrides,
  251. }
  252. }
  253. function runtime(fake: FakeSandbox, getSandbox: () => Promise<Sandbox> = async () => fake.sandbox): E2BSandboxService {
  254. return {
  255. cwd: '/workspace',
  256. runtimeRoot: '/workspace/.dsh-e2b',
  257. disposeMode: 'kill',
  258. getSandbox,
  259. } as unknown as E2BSandboxService
  260. }
  261. async function flush(): Promise<void> {
  262. await new Promise(resolve => setTimeout(resolve, 0))
  263. }
  264. describe('E2BOutputReader', () => {
  265. it('decodes base64 across arbitrary callback boundaries and rejects malformed framing', () => {
  266. const decoder = new E2BBase64Decoder()
  267. expect(decoder.push('')).toEqual(Buffer.alloc(0))
  268. expect(decoder.push('5')).toEqual(Buffer.alloc(0))
  269. expect(decoder.push('L2')).toEqual(Buffer.alloc(0))
  270. expect(decoder.push('g\n').toString()).toBe('你')
  271. expect(decoder.push('YQ==\nYg==\n').toString()).toBe('ab')
  272. expect(decoder.push(`${Buffer.from([0, 255]).toString('base64')}\n`)).toEqual(Buffer.from([0, 255]))
  273. expect(decoder.push(`${E2B_OUTPUT_COMPLETE_FRAME}\n`)).toEqual(Buffer.alloc(0))
  274. decoder.finish()
  275. expect(() => new E2BBase64Decoder().push('%\n')).toThrow('invalid base64')
  276. expect(() => new E2BBase64Decoder().push('AB==\n')).toThrow('invalid base64')
  277. expect(() => decoder.push(`${E2B_OUTPUT_COMPLETE_FRAME}\n`)).toThrow('duplicate output transport completion')
  278. expect(() => decoder.push('YQ==\n')).toThrow('continued after completion')
  279. const truncated = new E2BBase64Decoder()
  280. truncated.push('YQ')
  281. expect(() => { truncated.finish() }).toThrow('truncated base64')
  282. expect(() => { new E2BBase64Decoder().finish() }).toThrow('incomplete output transport')
  283. const interrupted = new E2BBase64Decoder()
  284. interrupted.push('YQ')
  285. expect(() => { interrupted.finish(false) }).not.toThrow()
  286. })
  287. it('keeps a byte-exact tail with independent whole-stream cursors', () => {
  288. const reader = new E2BOutputReader(4, 10, '/remote/spill')
  289. reader.push(Buffer.alloc(0))
  290. reader.push(Buffer.from('ab'))
  291. reader.push(Buffer.from('cdef'))
  292. expect(reader.size).toBe(6)
  293. expect(reader.readFrom(0)).toEqual({ text: 'cdef', nextOffset: 6, lossy: true, spillPath: '/remote/spill' })
  294. expect(reader.readFrom(2)).toEqual({ text: 'cdef', nextOffset: 6, lossy: false })
  295. expect(reader.readFrom(5)).toEqual({ text: 'f', nextOffset: 6, lossy: false })
  296. expect(reader.readFrom(99)).toEqual({ text: '', nextOffset: 6, lossy: false })
  297. reader.invalidateSpill()
  298. expect(reader.readFrom(0)).toEqual({ text: 'cdef', nextOffset: 6, lossy: true })
  299. })
  300. it('drops whole head chunks and withholds absent or over-cap spills', () => {
  301. const withoutSpill = new E2BOutputReader(2, undefined, '/unused')
  302. withoutSpill.push(Buffer.from('ab'))
  303. withoutSpill.push(Buffer.from('cd'))
  304. expect(withoutSpill.readFrom(0)).toEqual({ text: 'cd', nextOffset: 4, lossy: true })
  305. const overCap = new E2BOutputReader(2, 3, '/too-small')
  306. overCap.push(Buffer.from('abcd'))
  307. expect(overCap.readFrom(0)).toEqual({ text: 'cd', nextOffset: 4, lossy: true })
  308. expect(() => overCap.readFrom(-1)).toThrow(/non-negative safe integer/)
  309. expect(() => overCap.readFrom(1.5)).toThrow(/non-negative safe integer/)
  310. })
  311. })
  312. describe('E2BSubprocessHandle', () => {
  313. it('starts asynchronously, keeps secrets out of the command, and supports deferred piped stdin/output', async () => {
  314. const fake = new FakeSandbox()
  315. fake.processGroupId = '4343\n'
  316. fake.deferStart()
  317. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  318. argv: ['tool', 'argument with spaces'],
  319. stdio: { stdin: 'pipe', stdout: 'pipe', stderr: { maxBytes: 8, spill: { maxBytes: 32 } } },
  320. env: { PATH: '/bin', 'FOO-BAR': 'hyphen-value', DEEPSEEK_API_KEY: 'explicit-secret', DSH_MODE: 'test' },
  321. }), '/workspace/.dsh-e2b/processes/one')
  322. expect(handle.pid).toBe(-1)
  323. handle.stdin!.write('hello')
  324. handle.stdin!.end()
  325. fake.releaseStart()
  326. await flush()
  327. expect(handle.pid).toBe(4343)
  328. expect(fake.handle.sent.map(value => String(value))).toEqual(['hello'])
  329. expect(fake.handle.closes).toBe(1)
  330. expect(fake.startOptions?.envs).toBeUndefined()
  331. const command = fake.commandsSeen.find(value => value.includes('exec "$dsh_e2b_env_bin" -i'))!
  332. expect(command).toContain('"$dsh_e2b_setsid" --wait -- "$dsh_e2b_bash" -c')
  333. expect(command).not.toContain('DEEPSEEK_API_KEY')
  334. expect(command).not.toContain('DSH_MODE')
  335. expect(command).not.toContain('FOO-BAR')
  336. expect(command).not.toContain('explicit-secret')
  337. expect(command).not.toContain('hyphen-value')
  338. expect(command).not.toContain('${!dsh_e2b_name}')
  339. expect(fake.commandsSeen).toContain('env -0')
  340. expect(command).toContain('mapfile -d')
  341. expect(command).toContain('dsh_e2b_node="$(command -v node)"')
  342. expect(command).toContain('"$dsh_e2b_env_bin" -i "$dsh_e2b_node" -e')
  343. expect(command).toContain('exec "$dsh_e2b_env_bin" -i "${dsh_e2b_env[@]}"')
  344. expect(command).toContain('>&2 2>/dev/null')
  345. expect(command).not.toContain('2>/dev/null >&2')
  346. expect(command).toContain('base64')
  347. expect(fake.writtenFiles[0]).toEqual([
  348. '/workspace/.dsh-e2b/processes/one/pid',
  349. '/workspace/.dsh-e2b/processes/one/exit-code',
  350. '/workspace/.dsh-e2b/processes/one/environment',
  351. '/workspace/.dsh-e2b/processes/one/stderr.log',
  352. ])
  353. expect(fake.writtenFileData.get('/workspace/.dsh-e2b/processes/one/environment')).toBe(
  354. 'PATH=/bin\0KEEP=safe\0FOO-BAR=hyphen-value\0DEEPSEEK_API_KEY=explicit-secret\0DSH_MODE=test\0',
  355. )
  356. let piped = ''
  357. handle.stdout!.on('data', (chunk) => { piped += String(chunk) })
  358. await fake.stdout('pipe-data')
  359. await fake.stderr('err')
  360. fake.finish()
  361. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  362. expect(piped).toBe('pipe-data')
  363. expect(handle.collected.stderr!.readFrom(0)).toMatchObject({ text: 'err', lossy: false })
  364. expect(fake.removed).toContain('/workspace/.dsh-e2b/processes/one/stderr.log')
  365. await expect(handle.waitForExit()).resolves.toBe(true)
  366. })
  367. it('preserves UTF-8 bytes when the ASCII transport is split across callbacks', async () => {
  368. const fake = new FakeSandbox()
  369. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  370. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  371. }), '/runtime/split-utf8')
  372. await flush()
  373. const chunks: Buffer[] = []
  374. handle.stdout!.on('data', (chunk: Buffer) => { chunks.push(chunk) })
  375. for (const character of `${Buffer.from('A你好B').toString('base64')}\n`) {
  376. await fake.stdoutWire(character)
  377. }
  378. fake.finish()
  379. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  380. expect(Buffer.concat(chunks).toString('utf8')).toBe('A你好B')
  381. })
  382. it('rejects malformed output transport without confusing it with a consumer sink failure', async () => {
  383. const fake = new FakeSandbox()
  384. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/malformed-output')
  385. await flush()
  386. await fake.stdoutWire('%\n')
  387. fake.finish()
  388. await expect(handle.done).rejects.toThrow('invalid base64 output transport')
  389. const stderrFake = new FakeSandbox()
  390. const stderrHandle = new E2BSubprocessHandle(runtime(stderrFake), spec(), '/runtime/malformed-stderr')
  391. await flush()
  392. await stderrFake.stderrWire('%\n')
  393. stderrFake.finish()
  394. await expect(stderrHandle.done).rejects.toThrow('invalid base64 output transport')
  395. })
  396. it('rejects a naturally completed command whose encoder omits its completion frame', async () => {
  397. const fake = new FakeSandbox()
  398. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/incomplete-output')
  399. await flush()
  400. fake.alive = false
  401. fake.handle.succeed(0)
  402. await expect(handle.done).rejects.toThrow('incomplete output transport')
  403. })
  404. it('bounds descendant-held output draining and withholds the incomplete spill', async () => {
  405. const fake = new FakeSandbox()
  406. const handle = new E2BSubprocessHandle(runtime(fake), spec({ graceMs: 5 }), '/runtime/drain-bound')
  407. await flush()
  408. await fake.stdout('leader-output')
  409. fake.exitStatus = '0\n'
  410. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  411. expect(fake.handle.disconnects).toBe(1)
  412. expect(handle.collected.stdout?.readFrom(0)).toEqual({
  413. text: 'tput',
  414. nextOffset: 13,
  415. lossy: true,
  416. })
  417. expect(fake.removed).toContain('/runtime/drain-bound/stdout.log')
  418. handle.terminate()
  419. await expect(handle.waitForExit()).resolves.toBe(true)
  420. })
  421. it('accepts clean encoder completion inside the output-drain grace', async () => {
  422. const fake = new FakeSandbox()
  423. const handle = new E2BSubprocessHandle(runtime(fake), spec({ graceMs: 100 }), '/runtime/drain-complete')
  424. await flush()
  425. fake.exitStatus = '0\n'
  426. fake.afterStatusRead = () => {
  427. fake.afterStatusRead = undefined
  428. setTimeout(() => { fake.finish() }, 0)
  429. }
  430. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  431. expect(fake.handle.disconnects).toBe(0)
  432. })
  433. it('preserves a requested signal when output draining expires', async () => {
  434. const fake = new FakeSandbox()
  435. fake.trapsTerm = true
  436. fake.delaysKill = true
  437. fake.delaysKillCompletion = true
  438. fake.sdkKillStops = false
  439. const handle = new E2BSubprocessHandle(runtime(fake), spec({ graceMs: 5 }), '/runtime/drain-signal')
  440. await flush()
  441. handle.terminate()
  442. await vi.waitFor(() => { expect(fake.commandsSeen).toContain('kill -KILL -- -4242') })
  443. fake.exitStatus = '143\n'
  444. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  445. expect(fake.handle.disconnects).toBe(1)
  446. fake.alive = false
  447. await expect(handle.waitForExit()).resolves.toBe(true)
  448. })
  449. it('rejects an invalid direct-command exit status', async () => {
  450. const fake = new FakeSandbox()
  451. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/invalid-status')
  452. await flush()
  453. fake.exitStatus = '999\n'
  454. await expect(handle.done).rejects.toThrow('invalid exit code')
  455. handle.terminate()
  456. await expect(handle.waitForExit()).resolves.toBe(true)
  457. })
  458. it('surfaces deferred piped-stdin write and close failures as stream errors', async () => {
  459. const writeFake = new FakeSandbox()
  460. writeFake.deferStart()
  461. vi.spyOn(writeFake.handle, 'sendStdin').mockRejectedValueOnce('stdin rejected')
  462. const writeHandle = new E2BSubprocessHandle(runtime(writeFake), spec({
  463. stdio: { stdin: 'pipe', stdout: { maxBytes: 4 }, stderr: { maxBytes: 4 } },
  464. }), '/runtime/stdin-write-error')
  465. const writeError = once(writeHandle.stdin!, 'error')
  466. writeHandle.stdin!.write('input')
  467. writeFake.releaseStart()
  468. await expect(writeError).resolves.toMatchObject([{ message: 'stdin rejected' }])
  469. writeFake.finish()
  470. await writeHandle.done
  471. const closeFake = new FakeSandbox()
  472. vi.spyOn(closeFake.handle, 'closeStdin').mockRejectedValueOnce(new Error('close rejected'))
  473. const closeHandle = new E2BSubprocessHandle(runtime(closeFake), spec({
  474. stdio: { stdin: 'pipe', stdout: { maxBytes: 4 }, stderr: { maxBytes: 4 } },
  475. }), '/runtime/stdin-close-error')
  476. await flush()
  477. const closeError = once(closeHandle.stdin!, 'error')
  478. closeHandle.stdin!.end()
  479. await expect(closeError).resolves.toMatchObject([{ message: 'close rejected' }])
  480. closeFake.finish()
  481. await closeHandle.done
  482. })
  483. it('collects bounded tails, retains valid spills, and maps natural nonzero exits', async () => {
  484. const fake = new FakeSandbox()
  485. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  486. stdio: {
  487. stdin: { data: 'batch' },
  488. stdout: { maxBytes: 4, spill: { maxBytes: 16 } },
  489. stderr: { maxBytes: 3 },
  490. },
  491. }), '/runtime/two')
  492. await flush()
  493. await fake.stdout('abcdef')
  494. await fake.stderr('12345')
  495. fake.finish(7)
  496. await expect(handle.done).resolves.toEqual({ exitCode: 7, signal: null })
  497. expect(fake.handle.sent).toEqual(['batch'])
  498. expect(fake.handle.closes).toBe(1)
  499. expect(handle.collected.stdout!.readFrom(0)).toEqual({
  500. text: 'cdef',
  501. nextOffset: 6,
  502. lossy: true,
  503. spillPath: '/runtime/two/stdout.log',
  504. })
  505. expect(handle.collected.stderr!.readFrom(0)).toEqual({ text: '345', nextOffset: 5, lossy: true })
  506. expect(fake.removed).not.toContain('/runtime/two/stdout.log')
  507. })
  508. it('removes a spill once the complete stream exceeds its cap', async () => {
  509. const fake = new FakeSandbox()
  510. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  511. stdio: { stdin: 'ignore', stdout: { maxBytes: 2, spill: { maxBytes: 3 } }, stderr: 'inherit' },
  512. }), '/runtime/oversize')
  513. await flush()
  514. await fake.stdout('abcd')
  515. await fake.stderr('')
  516. fake.finish()
  517. await handle.done
  518. expect(handle.collected.stdout!.readFrom(0)).toEqual({ text: 'cd', nextOffset: 4, lossy: true })
  519. expect(fake.removed).toContain('/runtime/oversize/stdout.log')
  520. const command = fake.commandsSeen.find(value => value.includes('dsh_e2b_tee='))!
  521. expect(command).toContain('"$dsh_e2b_head" -c 3')
  522. expect(command).toContain('/runtime/oversize/stdout.log')
  523. expect(command).toContain('"$dsh_e2b_tee" --output-error=warn-nopipe')
  524. expect(command).not.toContain('tee -a')
  525. })
  526. it('contains remote spill-removal failures and routes empty inherited output', async () => {
  527. const fake = new FakeSandbox()
  528. fake.nextRemoveError = new Error('already removed')
  529. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  530. stdio: { stdin: 'ignore', stdout: 'inherit', stderr: { maxBytes: 4, spill: { maxBytes: 8 } } },
  531. }), '/runtime/remove-error')
  532. await flush()
  533. await fake.stdout('')
  534. await fake.stderr('')
  535. fake.finish()
  536. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  537. expect(fake.removed).toContain('/runtime/remove-error/stderr.log')
  538. })
  539. it('terminates a process group with TERM and reports the signal outcome', async () => {
  540. const fake = new FakeSandbox()
  541. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/term')
  542. await flush()
  543. handle.terminate()
  544. handle.terminate()
  545. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  546. await expect(handle.waitForExit()).resolves.toBe(true)
  547. expect(fake.commandsSeen).toContain('kill -TERM -- -4242')
  548. expect(fake.commandsSeen).not.toContain('kill -KILL -- -4242')
  549. const signals = fake.commandsSeen.filter(command => command.startsWith('kill -')).length
  550. fake.alive = true
  551. handle.terminate()
  552. await flush()
  553. expect(fake.alive).toBe(true)
  554. expect(fake.commandsSeen.filter(command => command.startsWith('kill -'))).toHaveLength(signals)
  555. })
  556. it('escalates a TERM-trapping process group to KILL and uses the SDK kill as fallback', async () => {
  557. const fake = new FakeSandbox()
  558. fake.trapsTerm = true
  559. fake.handle.killError = new Error('already gone')
  560. const handle = new E2BSubprocessHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/kill')
  561. await flush()
  562. handle.terminate()
  563. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  564. await expect(handle.waitForExit()).resolves.toBe(true)
  565. expect(fake.commandsSeen).toContain('kill -KILL -- -4242')
  566. expect(fake.handle.kills).toBe(1)
  567. })
  568. it('honors termination requested before asynchronous startup finishes', async () => {
  569. const fake = new FakeSandbox()
  570. fake.deferStart()
  571. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/deferred-kill')
  572. handle.terminate()
  573. fake.releaseStart()
  574. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  575. })
  576. it('kills through the provisional SDK handle before process-group publication', async () => {
  577. const fake = new FakeSandbox()
  578. fake.deferProcessGroupRead()
  579. fake.signalErrors.push(commandError(1), commandError(1))
  580. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/pre-publication-kill')
  581. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  582. handle.terminate()
  583. await vi.waitFor(() => { expect(fake.handle.kills).toBe(1) })
  584. expect(fake.alive).toBe(false)
  585. await expect(handle.waitForExit()).resolves.toBe(true)
  586. fake.releaseProcessGroupRead()
  587. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  588. })
  589. it('bounds a quiescence observer while provisional termination is awaiting the controller', async () => {
  590. const fake = new FakeSandbox()
  591. fake.deferProcessGroupRead()
  592. const reconnect = Promise.withResolvers<Sandbox>()
  593. let calls = 0
  594. const delayedRuntime = runtime(fake, async () => {
  595. calls += 1
  596. return calls === 1 ? fake.sandbox : await reconnect.promise
  597. })
  598. const handle = new E2BSubprocessHandle(delayedRuntime, spec(), '/runtime/pre-publication-observer')
  599. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  600. handle.terminate()
  601. const controller = new AbortController()
  602. const waiting = handle.waitForExit(controller.signal)
  603. await flush()
  604. controller.abort()
  605. await expect(waiting).resolves.toBe(false)
  606. reconnect.resolve(fake.sandbox)
  607. await expect(handle.waitForExit()).resolves.toBe(true)
  608. fake.releaseProcessGroupRead()
  609. await handle.done
  610. })
  611. it('proves a provisional group exit when the SDK kill fallback fails', async () => {
  612. const fake = new FakeSandbox()
  613. fake.deferProcessGroupRead()
  614. fake.trapsTerm = true
  615. fake.handle.killError = new Error('SDK kill unavailable')
  616. const handle = new E2BSubprocessHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/pre-publication-group-kill')
  617. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  618. handle.terminate()
  619. await expect(handle.waitForExit()).resolves.toBe(true)
  620. fake.releaseProcessGroupRead()
  621. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  622. })
  623. it('reports failed provisional group and SDK force transports', async () => {
  624. const fake = new FakeSandbox()
  625. fake.deferProcessGroupRead()
  626. fake.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  627. fake.handle.killError = new Error('SDK kill failed')
  628. const handle = new E2BSubprocessHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/pre-publication-failure')
  629. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  630. handle.terminate()
  631. await expect(handle.waitForExit()).rejects.toThrow('force termination failed through both')
  632. fake.handle.killError = undefined
  633. handle.terminate()
  634. await expect(handle.waitForExit()).resolves.toBe(true)
  635. fake.releaseProcessGroupRead()
  636. await handle.done
  637. const absentGroup = new FakeSandbox()
  638. absentGroup.deferProcessGroupRead()
  639. absentGroup.signalErrors.push(commandError(1), commandError(1))
  640. absentGroup.handle.killError = new Error('SDK kill failed without a provisional group')
  641. const absentHandle = new E2BSubprocessHandle(
  642. runtime(absentGroup),
  643. spec({ graceMs: 1 }),
  644. '/runtime/pre-publication-absent-group',
  645. )
  646. await vi.waitFor(() => { expect(absentGroup.startOptions).toBeDefined() })
  647. absentHandle.terminate()
  648. await expect(absentHandle.waitForExit()).rejects.toThrow('force termination failed through both')
  649. absentGroup.handle.killError = undefined
  650. absentHandle.terminate()
  651. await expect(absentHandle.waitForExit()).resolves.toBe(true)
  652. absentGroup.releaseProcessGroupRead()
  653. await absentHandle.done
  654. })
  655. it('honors an already-aborted signal when constructing the asynchronous handle directly', async () => {
  656. const fake = new FakeSandbox()
  657. const handle = new E2BSubprocessHandle(runtime(fake), spec({ signal: AbortSignal.abort('stop') }), '/runtime/pre-aborted')
  658. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  659. })
  660. it('reacts to a signal that aborts after the remote command has started', async () => {
  661. const fake = new FakeSandbox()
  662. const controller = new AbortController()
  663. const handle = new E2BSubprocessHandle(runtime(fake), spec({ signal: controller.signal }), '/runtime/live-abort')
  664. await flush()
  665. controller.abort('stop')
  666. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  667. })
  668. it('can terminate a surviving process group after the command leader settles', async () => {
  669. const fake = new FakeSandbox()
  670. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/surviving-group')
  671. await flush()
  672. await fake.completeOutput()
  673. fake.handle.succeed(0)
  674. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  675. expect(fake.alive).toBe(true)
  676. handle.terminate()
  677. await flush()
  678. const signaled = fake.commandsSeen.includes('kill -TERM -- -4242')
  679. if (!signaled) fake.finish()
  680. await expect(handle.waitForExit()).resolves.toBe(true)
  681. expect(signaled).toBe(true)
  682. })
  683. it('bounds waitForExit while startup or a live group is pending', async () => {
  684. const fake = new FakeSandbox()
  685. fake.deferStart()
  686. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/wait')
  687. const beforeStart = new AbortController()
  688. const pending = handle.waitForExit(beforeStart.signal)
  689. beforeStart.abort()
  690. await expect(pending).resolves.toBe(false)
  691. await expect(handle.waitForExit(AbortSignal.abort())).resolves.toBe(false)
  692. fake.releaseStart()
  693. await flush()
  694. const live = new AbortController()
  695. const liveWait = handle.waitForExit(live.signal)
  696. live.abort()
  697. await expect(liveWait).resolves.toBe(false)
  698. fake.finish()
  699. await handle.done
  700. const terminatingFake = new FakeSandbox()
  701. terminatingFake.deferStart()
  702. const terminating = new E2BSubprocessHandle(runtime(terminatingFake), spec(), '/runtime/wait-termination-start')
  703. terminating.terminate()
  704. const beforeHandle = new AbortController()
  705. const handlePending = terminating.waitForExit(beforeHandle.signal)
  706. beforeHandle.abort()
  707. await expect(handlePending).resolves.toBe(false)
  708. terminatingFake.releaseStart()
  709. await terminating.done
  710. })
  711. it('bounds both sides of the liveness-poll abort race', async () => {
  712. const fake = new FakeSandbox()
  713. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/poll-abort')
  714. await flush()
  715. const beforeTick = new AbortController()
  716. fake.afterProbe = () => { beforeTick.abort(); fake.afterProbe = undefined }
  717. await expect(handle.waitForExit(beforeTick.signal)).resolves.toBe(false)
  718. const duringTick = new AbortController()
  719. fake.afterProbe = () => {
  720. fake.afterProbe = undefined
  721. setTimeout(() => { duringTick.abort() }, 0)
  722. }
  723. await expect(handle.waitForExit(duringTick.signal)).resolves.toBe(false)
  724. const duringProbe = new AbortController()
  725. fake.beforeProbe = () => { duringProbe.abort(); fake.beforeProbe = undefined }
  726. await expect(handle.waitForExit(duringProbe.signal)).resolves.toBe(false)
  727. let racedAbort = false
  728. const raceSignal = {
  729. get aborted() { return racedAbort },
  730. addEventListener: () => { racedAbort = true },
  731. removeEventListener: () => {},
  732. } as unknown as AbortSignal
  733. await expect(handle.waitForExit(raceSignal)).resolves.toBe(false)
  734. fake.finish()
  735. await handle.done
  736. })
  737. it('observes a live group across one successful bounded poll', async () => {
  738. const fake = new FakeSandbox()
  739. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/poll-success')
  740. await flush()
  741. setTimeout(() => { fake.finish() }, 1)
  742. await expect(handle.waitForExit(new AbortController().signal)).resolves.toBe(true)
  743. await handle.done
  744. })
  745. it('treats startup failure as no live tree and contains readiness rejection', async () => {
  746. const fake = new FakeSandbox()
  747. fake.backgroundError = new Error('start failed')
  748. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/fail')
  749. await expect(handle.done).rejects.toThrow('start failed')
  750. expect(handle.pid).toBe(-1)
  751. expect(fake.removed).toContain('/runtime/fail/environment')
  752. expect(fake.removed).toContain('/runtime/fail')
  753. await expect(handle.waitForExit()).resolves.toBe(true)
  754. handle.terminate()
  755. const unavailableHandle = new E2BSubprocessHandle(
  756. runtime(new FakeSandbox(), async () => { throw new Error('sandbox unavailable') }),
  757. spec(),
  758. '/runtime/unavailable-start',
  759. )
  760. await expect(unavailableHandle.done).rejects.toThrow('sandbox unavailable')
  761. await expect(unavailableHandle.waitForExit()).resolves.toBe(true)
  762. const envFailure = new FakeSandbox()
  763. envFailure.envError = new Error('ambient lookup failed')
  764. const envHandle = new E2BSubprocessHandle(runtime(envFailure), spec(), '/runtime/env-failure')
  765. await expect(envHandle.done).rejects.toThrow('ambient lookup failed')
  766. expect(envFailure.removed).toEqual([])
  767. const cleanupFailure = new FakeSandbox()
  768. cleanupFailure.backgroundError = new Error('start failed before credential consumption')
  769. cleanupFailure.nextRemoveError = new Error('credential cleanup failed')
  770. const cleanupHandle = new E2BSubprocessHandle(runtime(cleanupFailure), spec(), '/runtime/cleanup-failure')
  771. await expect(cleanupHandle.done).rejects.toThrow('command failed and private state cleanup failed')
  772. const absentState = new FakeSandbox()
  773. absentState.backgroundError = new Error('start failed after external cleanup')
  774. absentState.nextRemoveError = new FileNotFoundError('already removed')
  775. const absentHandle = new E2BSubprocessHandle(runtime(absentState), spec(), '/runtime/absent-state')
  776. await expect(absentHandle.done).rejects.toThrow('start failed after external cleanup')
  777. })
  778. it('bounds a readiness rejection with a still-live caller signal', async () => {
  779. const fake = new FakeSandbox()
  780. fake.deferStart()
  781. fake.backgroundError = new Error('start failed')
  782. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/fail-with-signal')
  783. const waiting = handle.waitForExit(new AbortController().signal)
  784. fake.releaseStart()
  785. await expect(handle.done).rejects.toThrow('start failed')
  786. await expect(waiting).resolves.toBe(true)
  787. })
  788. it('propagates an unavailable sandbox unless the caller aborts the wait', async () => {
  789. const fake = new FakeSandbox()
  790. let calls = 0
  791. const unavailable = runtime(fake, async () => {
  792. calls += 1
  793. if (calls === 1) return fake.sandbox
  794. throw new Error('connection unavailable')
  795. })
  796. const handle = new E2BSubprocessHandle(unavailable, spec(), '/runtime/unavailable')
  797. await flush()
  798. await expect(handle.waitForExit()).rejects.toThrow('connection unavailable')
  799. fake.finish()
  800. await handle.done
  801. })
  802. it('returns false when the caller aborts while reconnecting for liveness', async () => {
  803. const fake = new FakeSandbox()
  804. const reconnect = Promise.withResolvers<Sandbox>()
  805. let calls = 0
  806. const unavailable = runtime(fake, async () => {
  807. calls += 1
  808. return calls === 1 ? fake.sandbox : await reconnect.promise
  809. })
  810. const handle = new E2BSubprocessHandle(unavailable, spec(), '/runtime/reconnect-abort')
  811. await flush()
  812. const controller = new AbortController()
  813. const waiting = handle.waitForExit(controller.signal)
  814. await flush()
  815. controller.abort()
  816. reconnect.reject(new Error('connection unavailable'))
  817. await expect(waiting).resolves.toBe(false)
  818. fake.finish()
  819. await handle.done
  820. })
  821. it('returns false when a liveness request itself is aborted and surfaces other probe failures', async () => {
  822. const fake = new FakeSandbox()
  823. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/probe')
  824. await flush()
  825. const controller = new AbortController()
  826. controller.abort()
  827. await expect(handle.waitForExit(controller.signal)).resolves.toBe(false)
  828. fake.probeError = new Error('probe failed')
  829. await expect(handle.waitForExit()).rejects.toThrow('probe failed')
  830. fake.finish()
  831. await handle.done
  832. })
  833. it('makes batch stdin close failures best-effort', async () => {
  834. const fake = new FakeSandbox()
  835. vi.spyOn(fake.handle, 'sendStdin').mockRejectedValueOnce(new Error('closed'))
  836. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  837. stdio: { stdin: { data: 'ignored' }, stdout: { maxBytes: 4 }, stderr: { maxBytes: 4 } },
  838. }), '/runtime/stdin-closed')
  839. await flush()
  840. fake.finish()
  841. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  842. })
  843. it('rejects malformed SDK process ids and non-command settlement failures', async () => {
  844. const invalidPid = new FakeSandbox()
  845. invalidPid.handle.pid = 0
  846. const invalid = new E2BSubprocessHandle(runtime(invalidPid), spec(), '/runtime/invalid-pid')
  847. await expect(invalid.done).rejects.toThrow(/invalid command pid 0/)
  848. expect(invalidPid.handle.kills).toBe(1)
  849. expect(invalidPid.removed).toContain('/runtime/invalid-pid/environment')
  850. await expect(invalid.waitForExit()).resolves.toBe(true)
  851. const failedRollback = new FakeSandbox()
  852. failedRollback.handle.pid = 0
  853. failedRollback.handle.killError = new Error('invalid handle kill failed')
  854. const retained = new E2BSubprocessHandle(runtime(failedRollback), spec(), '/runtime/invalid-pid-retained')
  855. await expect(retained.done).rejects.toThrow('invalid command pid rollback did not reach quiescence')
  856. await expect(retained.waitForExit()).rejects.toThrow('invalid handle kill failed')
  857. failedRollback.handle.killError = undefined
  858. retained.terminate()
  859. await expect(retained.waitForExit()).resolves.toBe(true)
  860. const crashedFake = new FakeSandbox()
  861. const crashed = new E2BSubprocessHandle(runtime(crashedFake), spec(), '/runtime/crashed')
  862. await flush()
  863. crashedFake.alive = false
  864. crashedFake.handle.crash(new Error('command transport failed'))
  865. await expect(crashed.done).rejects.toThrow('command transport failed')
  866. })
  867. it('rejects invalid or absent process-group publication', async () => {
  868. const invalidGroup = new FakeSandbox()
  869. invalidGroup.processGroupId = 'not-a-pid\n'
  870. invalidGroup.delaysKill = true
  871. invalidGroup.sdkKillStops = false
  872. invalidGroup.afterProbe = () => { invalidGroup.alive = false }
  873. const invalid = new E2BSubprocessHandle(runtime(invalidGroup), spec(), '/runtime/invalid-group')
  874. await expect(invalid.done).rejects.toThrow(/invalid process-group id/)
  875. expect(invalidGroup.handle.kills).toBe(1)
  876. expect(invalidGroup.commandsSeen).toContain('kill -KILL -- -4242')
  877. await expect(invalid.waitForExit()).resolves.toBe(true)
  878. const absentGroup = new FakeSandbox()
  879. absentGroup.processGroupId = ''
  880. const absent = new E2BSubprocessHandle(runtime(absentGroup), spec(), '/runtime/absent-group')
  881. await flush()
  882. absentGroup.finish()
  883. await expect(absent.done).rejects.toThrow(/exited before publishing/)
  884. expect(absentGroup.handle.kills).toBe(1)
  885. expect(absentGroup.commandsSeen).toContain('kill -KILL -- -4242')
  886. await expect(absent.waitForExit()).resolves.toBe(true)
  887. })
  888. it('preserves publication and rollback failures when cleanup cannot be verified', async () => {
  889. const fake = new FakeSandbox()
  890. fake.processGroupId = 'not-a-pid\n'
  891. fake.signalError = new Error('rollback signal failed')
  892. fake.handle.killError = new Error('SDK kill failed')
  893. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/failed-rollback')
  894. let failure: unknown
  895. try {
  896. await handle.done
  897. } catch (error: unknown) {
  898. failure = error
  899. }
  900. expect(failure).toBeInstanceOf(AggregateError)
  901. if (!(failure instanceof AggregateError)) throw new Error('expected AggregateError')
  902. expect(failure.message).toBe('subprocess-e2b: process-group publication failed and rollback did not reach quiescence')
  903. const failures = Array.from(failure.errors as Iterable<unknown>)
  904. expect(failures).toHaveLength(2)
  905. expect(failures[0]).toBeInstanceOf(Error)
  906. expect(failures[1]).toBeInstanceOf(Error)
  907. if (!(failures[0] instanceof Error) || !(failures[1] instanceof Error)) throw new Error('expected nested errors')
  908. expect(failures[0].message).toContain('invalid process-group id')
  909. expect(failures[1].message).toBe('rollback signal failed')
  910. expect(fake.handle.kills).toBe(1)
  911. const bounded = new AbortController()
  912. const waiting = handle.waitForExit(bounded.signal)
  913. bounded.abort()
  914. await expect(waiting).resolves.toBe(false)
  915. handle.terminate()
  916. await expect(handle.waitForExit()).resolves.toBe(true)
  917. expect(fake.commandsSeen).toContain('kill -TERM -- -4242')
  918. })
  919. it('waits for delayed process-group publication', async () => {
  920. const fake = new FakeSandbox()
  921. fake.processGroupReads.push('', '4242\n')
  922. const handle = new E2BSubprocessHandle(runtime(fake), spec(), '/runtime/delayed-group')
  923. await vi.waitFor(() => { expect(handle.pid).toBe(4242) })
  924. fake.finish()
  925. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  926. })
  927. it('handles output backpressure and contains a stderr sink failure', async () => {
  928. const fake = new FakeSandbox()
  929. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  930. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: 'pipe' },
  931. }), '/runtime/backpressure')
  932. await flush()
  933. handle.stdout!.on('error', () => {})
  934. const stdoutWrite = vi.spyOn(handle.stdout!, 'write').mockReturnValueOnce(false)
  935. const stdoutPending = fake.stdout('blocked')
  936. queueMicrotask(() => { handle.stdout!.emit('drain') })
  937. await stdoutPending
  938. stdoutWrite.mockRestore()
  939. handle.stderr!.on('error', () => {})
  940. const stderrWrite = vi.spyOn(handle.stderr!, 'write').mockReturnValueOnce(false)
  941. const stderrPending = fake.stderr('broken')
  942. queueMicrotask(() => { handle.stderr!.emit('error', new Error('sink failed')) })
  943. await stderrPending
  944. stderrWrite.mockRestore()
  945. fake.finish()
  946. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  947. })
  948. it('contains a pipe callback failure instead of rejecting command settlement', async () => {
  949. const fake = new FakeSandbox()
  950. const handle = new E2BSubprocessHandle(runtime(fake), spec({
  951. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  952. }), '/runtime/pipe-error')
  953. await flush()
  954. const emitted = once(handle.stdout!, 'error')
  955. handle.stdout!.destroy(new Error('consumer failed'))
  956. await emitted
  957. await fake.stdout('late output')
  958. fake.finish()
  959. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  960. })
  961. it('contains an already-gone group signal and escalates after a TERM transport failure', async () => {
  962. const gone = new FakeSandbox()
  963. gone.trapsTerm = true
  964. gone.signalError = commandError(1)
  965. const goneHandle = new E2BSubprocessHandle(runtime(gone), spec({ graceMs: 1 }), '/runtime/gone-signal')
  966. await flush()
  967. goneHandle.terminate()
  968. await expect(goneHandle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  969. const failed = new FakeSandbox()
  970. failed.signalError = new Error('signal transport failed')
  971. const failedHandle = new E2BSubprocessHandle(runtime(failed), spec(), '/runtime/failed-signal')
  972. await flush()
  973. failedHandle.terminate()
  974. await expect(failedHandle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  975. expect(failed.commandsSeen).toContain('kill -KILL -- -4242')
  976. })
  977. it('allows termination retry after both force transports fail', async () => {
  978. const fake = new FakeSandbox()
  979. fake.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  980. fake.handle.killError = new Error('SDK kill failed')
  981. const handle = new E2BSubprocessHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/retry-signal')
  982. await flush()
  983. handle.terminate()
  984. await expect(handle.waitForExit()).rejects.toThrow('force termination failed through both')
  985. fake.handle.killError = undefined
  986. handle.terminate()
  987. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  988. expect(fake.commandsSeen.filter(command => command.startsWith('kill -TERM '))).toHaveLength(2)
  989. const missingGroup = new FakeSandbox()
  990. missingGroup.trapsTerm = true
  991. missingGroup.signalErrors.push(undefined, commandError(1))
  992. missingGroup.handle.killError = new Error('SDK kill failed after group exit race')
  993. const raced = new E2BSubprocessHandle(runtime(missingGroup), spec({ graceMs: 1 }), '/runtime/group-exit-race')
  994. await flush()
  995. raced.terminate()
  996. await expect(raced.waitForExit()).rejects.toThrow('force termination failed through both')
  997. missingGroup.handle.killError = undefined
  998. raced.terminate()
  999. await expect(raced.waitForExit()).resolves.toBe(true)
  1000. })
  1001. })
  1002. describe('E2BSubprocessService', () => {
  1003. async function service(
  1004. fake = new FakeSandbox(),
  1005. providedRuntime: E2BSandboxService = runtime(fake),
  1006. ): Promise<{ ctx: Context; fiber: Awaited<ReturnType<Context['plugin']>> }> {
  1007. const ctx = new Context()
  1008. ctx.provide('e2b', providedRuntime)
  1009. const fiber = await ctx.plugin(E2BSubprocessService)
  1010. return { ctx, fiber }
  1011. }
  1012. it('registers handles and disposal terminates and joins live remote groups regardless of sandbox policy', async () => {
  1013. const fake = new FakeSandbox()
  1014. fake.trapsTerm = true
  1015. const { ctx, fiber } = await service(fake)
  1016. const handle = ctx.subprocess.spawn(spec({ graceMs: 1 }))
  1017. await flush()
  1018. await fiber.dispose()
  1019. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  1020. expect(fake.alive).toBe(false)
  1021. })
  1022. it('awaits SDK settlement after the remote process group becomes quiescent', async () => {
  1023. const fake = new FakeSandbox()
  1024. fake.trapsTerm = true
  1025. const { ctx, fiber } = await service(fake)
  1026. const handle = ctx.subprocess.spawn(spec())
  1027. await flush()
  1028. fake.alive = false
  1029. let disposed = false
  1030. const disposing = fiber.dispose().then(() => { disposed = true })
  1031. await flush()
  1032. expect(disposed).toBe(false)
  1033. fake.finish()
  1034. await disposing
  1035. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  1036. })
  1037. it('reports a failed termination transaction from disposal instead of waiting on done', async () => {
  1038. const fake = new FakeSandbox()
  1039. fake.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  1040. fake.handle.killError = new Error('SDK kill failed')
  1041. const { ctx, fiber } = await service(fake)
  1042. const handle = ctx.subprocess.spawn(spec({ graceMs: 1 }))
  1043. await flush()
  1044. await expect(fiber.dispose()).resolves.toBeUndefined()
  1045. await expect(handle.waitForExit()).rejects.toThrow('force termination failed through both')
  1046. fake.handle.killError = undefined
  1047. handle.terminate()
  1048. await expect(handle.waitForExit()).resolves.toBe(true)
  1049. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  1050. })
  1051. it('releases naturally settled handles before later service disposal', async () => {
  1052. const fake = new FakeSandbox()
  1053. const { ctx, fiber } = await service(fake)
  1054. const handle = ctx.subprocess.spawn(spec())
  1055. await flush()
  1056. fake.finish()
  1057. await handle.done
  1058. await flush()
  1059. const signalsBefore = fake.commandsSeen.filter(command => command.startsWith('kill -')).length
  1060. await fiber.dispose()
  1061. expect(fake.commandsSeen.filter(command => command.startsWith('kill -')).length).toBe(signalsBefore)
  1062. })
  1063. it('contains a release liveness failure and retries quiescence during disposal', async () => {
  1064. const fake = new FakeSandbox()
  1065. let calls = 0
  1066. const reconnecting = runtime(fake, async () => {
  1067. calls += 1
  1068. if (calls === 2) throw new Error('transient liveness failure')
  1069. return fake.sandbox
  1070. })
  1071. const { ctx, fiber } = await service(fake, reconnecting)
  1072. const handle = ctx.subprocess.spawn(spec())
  1073. await flush()
  1074. fake.finish()
  1075. await handle.done
  1076. await flush()
  1077. await fiber.dispose()
  1078. expect(calls).toBeGreaterThanOrEqual(3)
  1079. })
  1080. it('contains spawn rejection while disposal is joining the pending handle', async () => {
  1081. const fake = new FakeSandbox()
  1082. fake.deferStart()
  1083. fake.backgroundError = new Error('start failed during disposal')
  1084. const { ctx, fiber } = await service(fake)
  1085. const subprocess = ctx.subprocess
  1086. const handle = subprocess.spawn(spec())
  1087. const disposing = fiber.dispose()
  1088. await flush()
  1089. expect(() => subprocess.spawn(spec())).toThrow('service is disposing')
  1090. fake.releaseStart()
  1091. await expect(disposing).resolves.toBeUndefined()
  1092. await expect(handle.done).rejects.toThrow('start failed during disposal')
  1093. })
  1094. it('validates synchronous spawn preconditions', async () => {
  1095. const { ctx } = await service()
  1096. expect(() => ctx.subprocess.spawn(spec({ argv: [] }))).toThrow(/non-empty program/)
  1097. expect(() => ctx.subprocess.spawn(spec({ graceMs: 0 }))).toThrow(/positive finite/)
  1098. expect(() => ctx.subprocess.spawn(spec({ signal: AbortSignal.abort('stop') }))).toThrow(/aborted before spawn/)
  1099. expect(() => ctx.subprocess.spawn(spec({ signal: { aborted: true, reason: undefined } as AbortSignal }))).toThrow(/aborted$/)
  1100. })
  1101. it('registers the package-owned empty invariant installer', async () => {
  1102. const ctx = new Context()
  1103. await ctx.plugin(InvariantService, { enabled: true })
  1104. const fiber = await ctx.plugin(E2BSubprocessInvariant).await()
  1105. await fiber.dispose()
  1106. })
  1107. })