subprocess.spec.ts 71 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791
  1. import { once } from 'node:events'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import {
  4. CommandExitError,
  5. FileNotFoundError,
  6. SandboxNotFoundError,
  7. type CommandHandle,
  8. type CommandResult,
  9. type Sandbox,
  10. } from '@deepseek-ai/dsh-e2b'
  11. import type E2BRuntime from '@deepseek-ai/dsh-e2b'
  12. import type { SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'
  13. import E2BSubprocessRuntime from '@deepseek-ai/dsh-subprocess-e2b'
  14. import * as E2BSubprocessInvariant from '../src/invariant.ts'
  15. import { E2BBase64Decoder, E2B_OUTPUT_COMPLETE_FRAME, E2BOutputReader } from '../src/output.ts'
  16. import { E2BSubprocessHandle } from '../src/process.ts'
  17. import InvariantRegistry from '@deepseek-ai/dsh-invariants'
  18. import { describe, expect, it, vi } from 'vitest'
  19. function commandError(exitCode: number): CommandExitError {
  20. return new CommandExitError({ exitCode, stdout: '', stderr: '', error: `exit ${exitCode}` })
  21. }
  22. interface StartOptions {
  23. background: true
  24. cwd: string
  25. stdin: boolean
  26. timeoutMs: number
  27. signal?: AbortSignal
  28. envs?: Record<string, string>
  29. onStdout?: (data: string) => void | Promise<void>
  30. onStderr?: (data: string) => void | Promise<void>
  31. }
  32. class FakeCommandHandle {
  33. pid = 4242
  34. readonly sent: Array<string | Uint8Array> = []
  35. closes = 0
  36. kills = 0
  37. disconnects = 0
  38. killError: unknown
  39. killResult = true
  40. disconnectError: unknown
  41. private readonly result = Promise.withResolvers<CommandResult>()
  42. private settled = false
  43. constructor(private readonly onKill: () => void = () => {}) {}
  44. wait(): Promise<CommandResult> {
  45. return this.result.promise
  46. }
  47. async sendStdin(data: string | Uint8Array): Promise<void> {
  48. this.sent.push(data)
  49. }
  50. async closeStdin(): Promise<void> {
  51. this.closes += 1
  52. }
  53. async kill(): Promise<boolean> {
  54. this.kills += 1
  55. if (this.killError !== undefined) throw this.killError
  56. this.onKill()
  57. return this.killResult
  58. }
  59. async disconnect(): Promise<void> {
  60. this.disconnects += 1
  61. if (this.disconnectError !== undefined) throw this.disconnectError
  62. }
  63. succeed(exitCode = 0): void {
  64. if (this.settled) return
  65. this.settled = true
  66. this.result.resolve({ exitCode, stdout: '', stderr: '' })
  67. }
  68. fail(exitCode: number): void {
  69. if (this.settled) return
  70. this.settled = true
  71. this.result.reject(commandError(exitCode))
  72. }
  73. crash(error: unknown): void {
  74. if (this.settled) return
  75. this.settled = true
  76. this.result.reject(error)
  77. }
  78. }
  79. class FakeSandbox {
  80. readonly handle: FakeCommandHandle
  81. readonly commandsSeen: string[] = []
  82. readonly writtenFiles: string[][] = []
  83. readonly writtenFileData = new Map<string, string>()
  84. readonly removed: string[] = []
  85. readonly directories: string[] = []
  86. startOptions: StartOptions | undefined
  87. backgroundError: unknown
  88. envError: unknown
  89. statusError: unknown
  90. nextRemoveError: unknown
  91. probeError: unknown
  92. signalError: unknown
  93. readonly signalErrors: unknown[] = []
  94. trapsTerm = false
  95. delaysKill = false
  96. delaysKillCompletion = false
  97. sdkKillStops = true
  98. alive = true
  99. zombieOnly = false
  100. ambient = 'PATH=/ambient/bin\0KEEP=safe\0UNICODE=你好\0NPM_TOKEN=secret\0DSH_STALE=old\0BROKEN\0=bad\0'
  101. environmentHome = '/home/user'
  102. environmentWire: string | undefined
  103. environmentRequest: ((signal: AbortSignal | undefined) => Promise<void>) | undefined
  104. processGroupId = '4242\n'
  105. exitStatus = ''
  106. statusReads = 0
  107. readonly processGroupReads: string[] = []
  108. afterStatusRead: (() => void) | undefined
  109. beforeProbe: (() => void) | undefined
  110. afterProbe: (() => void) | undefined
  111. private startGate: Promise<void> | undefined
  112. private openStart: (() => void) | undefined
  113. private processGroupReadGate: Promise<void> | undefined
  114. private openProcessGroupRead: (() => void) | undefined
  115. private signalGate: Promise<void> | undefined
  116. private openSignal: (() => void) | undefined
  117. constructor() {
  118. this.handle = new FakeCommandHandle(() => {
  119. if (this.sdkKillStops) {
  120. this.alive = false
  121. this.handle.fail(137)
  122. }
  123. })
  124. }
  125. deferStart(): void {
  126. const gate = Promise.withResolvers<undefined>()
  127. this.startGate = gate.promise
  128. this.openStart = () => { gate.resolve(undefined) }
  129. }
  130. releaseStart(): void {
  131. this.openStart?.()
  132. }
  133. deferProcessGroupRead(): void {
  134. const gate = Promise.withResolvers<undefined>()
  135. this.processGroupReadGate = gate.promise
  136. this.openProcessGroupRead = () => { gate.resolve(undefined) }
  137. }
  138. releaseProcessGroupRead(): void {
  139. this.openProcessGroupRead?.()
  140. }
  141. deferSignals(): void {
  142. const gate = Promise.withResolvers<undefined>()
  143. this.signalGate = gate.promise
  144. this.openSignal = () => { gate.resolve(undefined) }
  145. }
  146. releaseSignals(): void {
  147. this.openSignal?.()
  148. }
  149. finish(exitCode = 0): void {
  150. this.alive = false
  151. void this.completeOutput().then(
  152. () => {
  153. if (exitCode === 0) this.handle.succeed(0)
  154. else this.handle.fail(exitCode)
  155. },
  156. (error: unknown) => { this.handle.crash(error) },
  157. )
  158. }
  159. async completeOutput(): Promise<void> {
  160. await Promise.all([
  161. this.stdoutWire(`${E2B_OUTPUT_COMPLETE_FRAME}\n`),
  162. this.stderrWire(`${E2B_OUTPUT_COMPLETE_FRAME}\n`),
  163. ])
  164. }
  165. async stdout(data: string): Promise<void> {
  166. await this.stdoutWire(data.length === 0 ? '' : `${Buffer.from(data).toString('base64')}\n`)
  167. }
  168. async stderr(data: string): Promise<void> {
  169. await this.stderrWire(data.length === 0 ? '' : `${Buffer.from(data).toString('base64')}\n`)
  170. }
  171. async stdoutWire(data: string): Promise<void> {
  172. await this.startOptions?.onStdout?.(data)
  173. }
  174. async stderrWire(data: string): Promise<void> {
  175. await this.startOptions?.onStderr?.(data)
  176. }
  177. readonly sandbox = {
  178. sandboxId: 'fake',
  179. files: {
  180. makeDir: async (path: string): Promise<boolean> => {
  181. this.directories.push(path)
  182. return true
  183. },
  184. write: async (files: Array<{ path: string; data: string }>): Promise<object[]> => {
  185. this.writtenFiles.push(files.map(file => file.path))
  186. for (const file of files) this.writtenFileData.set(file.path, file.data)
  187. return files.map(() => ({}))
  188. },
  189. read: async (path: string): Promise<string> => {
  190. if (!path.endsWith('/exit-code')) {
  191. await this.processGroupReadGate
  192. return this.processGroupReads.shift() ?? this.processGroupId
  193. }
  194. if (this.statusError !== undefined) {
  195. const error = this.statusError
  196. this.statusError = undefined
  197. throw error
  198. }
  199. this.statusReads += 1
  200. this.afterStatusRead?.()
  201. return this.exitStatus
  202. },
  203. remove: async (path: string): Promise<void> => {
  204. this.removed.push(path)
  205. if (this.nextRemoveError !== undefined) {
  206. const error = this.nextRemoveError
  207. this.nextRemoveError = undefined
  208. throw error
  209. }
  210. },
  211. },
  212. commands: {
  213. run: async (command: string, options?: StartOptions | { signal?: AbortSignal }): Promise<CommandHandle | CommandResult> => {
  214. this.commandsSeen.push(command)
  215. if (command.includes('env -0 | base64')) {
  216. await this.environmentRequest?.(options?.signal)
  217. if (this.envError !== undefined) throw this.envError
  218. return {
  219. exitCode: 0,
  220. stdout: this.environmentWire ?? [this.environmentHome, this.ambient]
  221. .map(value => Buffer.from(value).toString('base64'))
  222. .join('\n'),
  223. stderr: '',
  224. }
  225. }
  226. if (command.startsWith('set -o pipefail; ps -eo pgid=,stat=')) {
  227. this.beforeProbe?.()
  228. if (options?.signal?.aborted === true) throw new DOMException('aborted', 'AbortError')
  229. if (this.probeError !== undefined) {
  230. const error = this.probeError
  231. this.probeError = undefined
  232. throw error
  233. }
  234. const stdout = this.alive && !this.zombieOnly ? 'live\n' : ''
  235. this.afterProbe?.()
  236. return { exitCode: 0, stdout, stderr: '' }
  237. }
  238. if (command.startsWith('kill -TERM ')) {
  239. await this.signalGate
  240. const error = this.signalErrors.shift() ?? this.signalError
  241. if (error !== undefined) {
  242. if (this.signalErrors.length === 0) this.signalError = undefined
  243. throw error
  244. }
  245. if (!this.trapsTerm) {
  246. this.alive = false
  247. this.handle.fail(143)
  248. }
  249. return { exitCode: 0, stdout: '', stderr: '' }
  250. }
  251. if (command.startsWith('kill -KILL ')) {
  252. await this.signalGate
  253. const error = this.signalErrors.shift() ?? this.signalError
  254. if (error !== undefined) {
  255. if (this.signalErrors.length === 0) this.signalError = undefined
  256. throw error
  257. }
  258. if (!this.delaysKill) this.alive = false
  259. if (!this.delaysKillCompletion) this.handle.fail(137)
  260. return { exitCode: 0, stdout: '', stderr: '' }
  261. }
  262. if ((options as StartOptions | undefined)?.background === true) {
  263. this.startOptions = options as StartOptions
  264. await this.startGate
  265. if (this.backgroundError !== undefined) throw this.backgroundError
  266. return this.handle as unknown as CommandHandle
  267. }
  268. return { exitCode: 0, stdout: '', stderr: '' }
  269. },
  270. },
  271. } as unknown as Sandbox
  272. }
  273. function spec(overrides: Partial<SubprocessSpawnSpec> = {}): SubprocessSpawnSpec {
  274. return {
  275. argv: ['bash', '-c', 'printf ok'],
  276. cwd: '/workspace',
  277. stdio: {
  278. stdin: 'ignore',
  279. stdout: { maxBytes: 4, spill: { maxBytes: 16 } },
  280. stderr: { maxBytes: 4 },
  281. },
  282. graceMs: 5,
  283. ...overrides,
  284. }
  285. }
  286. function runtime(fake: FakeSandbox, getSandbox: () => Promise<Sandbox> = async () => fake.sandbox): E2BRuntime {
  287. return {
  288. cwd: '/workspace',
  289. runtimeRoot: '/workspace/.dsh-e2b',
  290. getSandbox,
  291. } as unknown as E2BRuntime
  292. }
  293. async function flush(): Promise<void> {
  294. await new Promise(resolve => setTimeout(resolve, 0))
  295. }
  296. /** Construct the handle under test with the config default the service would pass. */
  297. function testHandle(
  298. runtime: ConstructorParameters<typeof E2BSubprocessHandle>[0],
  299. spec: ConstructorParameters<typeof E2BSubprocessHandle>[1],
  300. stateDir: string,
  301. pollMs = 20,
  302. ): E2BSubprocessHandle {
  303. return new E2BSubprocessHandle(runtime, spec, stateDir, pollMs)
  304. }
  305. describe('E2BOutputReader', () => {
  306. it('decodes base64 across arbitrary callback boundaries and rejects malformed framing', () => {
  307. const decoder = new E2BBase64Decoder()
  308. expect(decoder.push('')).toEqual(Buffer.alloc(0))
  309. expect(decoder.push('5')).toEqual(Buffer.alloc(0))
  310. expect(decoder.push('L2')).toEqual(Buffer.alloc(0))
  311. expect(decoder.push('g\n').toString()).toBe('你')
  312. expect(decoder.push('YQ==\nYg==\n').toString()).toBe('ab')
  313. expect(decoder.push(`${Buffer.from([0, 255]).toString('base64')}\n`)).toEqual(Buffer.from([0, 255]))
  314. expect(decoder.push(`${E2B_OUTPUT_COMPLETE_FRAME}\n`)).toEqual(Buffer.alloc(0))
  315. decoder.finish()
  316. expect(() => new E2BBase64Decoder().push('%\n')).toThrow('invalid base64')
  317. expect(() => new E2BBase64Decoder().push('AB==\n')).toThrow('invalid base64')
  318. expect(() => decoder.push(`${E2B_OUTPUT_COMPLETE_FRAME}\n`)).toThrow('duplicate output transport completion')
  319. expect(() => decoder.push('YQ==\n')).toThrow('continued after completion')
  320. const truncated = new E2BBase64Decoder()
  321. truncated.push('YQ')
  322. expect(() => { truncated.finish() }).toThrow('truncated base64')
  323. expect(() => { new E2BBase64Decoder().finish() }).toThrow('incomplete output transport')
  324. const interrupted = new E2BBase64Decoder()
  325. interrupted.push('YQ')
  326. expect(() => { interrupted.finish(false) }).not.toThrow()
  327. })
  328. it('keeps a byte-exact tail with independent whole-stream cursors', () => {
  329. const reader = new E2BOutputReader(4, 10, '/remote/spill')
  330. reader.push(Buffer.alloc(0))
  331. reader.push(Buffer.from('ab'))
  332. reader.push(Buffer.from('cdef'))
  333. expect(reader.size).toBe(6)
  334. expect(reader.readFrom(0)).toEqual({ text: 'cdef', nextOffset: 6, lossy: true, spillPath: '/remote/spill' })
  335. expect(reader.readFrom(2)).toEqual({ text: 'cdef', nextOffset: 6, lossy: false })
  336. expect(reader.readFrom(5)).toEqual({ text: 'f', nextOffset: 6, lossy: false })
  337. expect(reader.readFrom(99)).toEqual({ text: '', nextOffset: 6, lossy: false })
  338. reader.invalidateSpill()
  339. expect(reader.readFrom(0)).toEqual({ text: 'cdef', nextOffset: 6, lossy: true })
  340. })
  341. it('drops whole head chunks and withholds absent or over-cap spills', () => {
  342. const withoutSpill = new E2BOutputReader(2, undefined, '/unused')
  343. withoutSpill.push(Buffer.from('ab'))
  344. withoutSpill.push(Buffer.from('cd'))
  345. expect(withoutSpill.readFrom(0)).toEqual({ text: 'cd', nextOffset: 4, lossy: true })
  346. const overCap = new E2BOutputReader(2, 3, '/too-small')
  347. overCap.push(Buffer.from('abcd'))
  348. expect(overCap.readFrom(0)).toEqual({ text: 'cd', nextOffset: 4, lossy: true })
  349. })
  350. })
  351. describe('E2BSubprocessHandle', () => {
  352. it('starts asynchronously, keeps secrets out of the command, and supports deferred piped stdin/output', async () => {
  353. const fake = new FakeSandbox()
  354. fake.processGroupId = '4343\n'
  355. fake.deferStart()
  356. const handle = testHandle(runtime(fake), spec({
  357. argv: ['tool', 'argument with spaces'],
  358. stdio: { stdin: 'pipe', stdout: 'pipe', stderr: { maxBytes: 8, spill: { maxBytes: 32 } } },
  359. env: {
  360. PATH: '/bin',
  361. 'FOO-BAR': 'hyphen-value',
  362. '--split-string': 'literal-value',
  363. DEEPSEEK_API_KEY: 'explicit-secret',
  364. DSH_MODE: 'test',
  365. // The seam's tombstone: an explicit undefined removes the ambient entry.
  366. KEEP: undefined,
  367. },
  368. }), '/workspace/.dsh-e2b/processes/one')
  369. handle.stdin!.write('hello')
  370. handle.stdin!.end()
  371. fake.releaseStart()
  372. await flush()
  373. expect(fake.handle.sent.map(value => String(value))).toEqual(['hello'])
  374. expect(fake.handle.closes).toBe(1)
  375. const controlEnvs = fake.startOptions?.envs
  376. expect(controlEnvs?.HOME).toMatch(/^\/\.dsh-e2b-control-/)
  377. expect(controlEnvs).toEqual({
  378. TERM: 'dumb',
  379. NPM_TOKEN: '',
  380. DSH_STALE: '',
  381. HOME: controlEnvs?.HOME,
  382. })
  383. const command = fake.commandsSeen.find(value => value.includes('exec "$dsh_e2b_env_bin" -i'))!
  384. expect(command).toContain('"$dsh_e2b_setsid" --wait -- "$dsh_e2b_bash" -c')
  385. expect(command).not.toContain('DEEPSEEK_API_KEY')
  386. expect(command).not.toContain('DSH_MODE')
  387. expect(command).not.toContain('FOO-BAR')
  388. expect(command).not.toContain('explicit-secret')
  389. expect(command).not.toContain('hyphen-value')
  390. expect(command).not.toContain('${!dsh_e2b_name}')
  391. const environmentProbe = fake.commandsSeen.find(value => value.includes('env -0 | base64'))
  392. expect(environmentProbe).toContain('getent passwd "$(id -u)"')
  393. expect(environmentProbe).toContain('test -n "$dsh_e2b_home" -a -d "$dsh_e2b_home"')
  394. expect(environmentProbe).not.toContain('"$PWD"')
  395. expect(command).toContain('mapfile -d')
  396. expect(command).toContain('dsh_e2b_node="$(command -v node)"')
  397. expect(command).toContain('"$dsh_e2b_env_bin" -i "$dsh_e2b_node" -e')
  398. expect(command).toContain('"$dsh_e2b_env_bin" -i -- "${dsh_e2b_env[@]}" "$@"')
  399. expect(command).toContain('exec "$dsh_e2b_env_bin" -i -- "${dsh_e2b_env[@]}"')
  400. expect(command).toContain('>&2 2>/dev/null')
  401. expect(command).not.toContain('2>/dev/null >&2')
  402. expect(command).toContain('base64')
  403. expect(fake.writtenFiles[0]).toEqual([
  404. '/workspace/.dsh-e2b/processes/one/pid',
  405. '/workspace/.dsh-e2b/processes/one/exit-code',
  406. '/workspace/.dsh-e2b/processes/one/environment',
  407. '/workspace/.dsh-e2b/processes/one/stderr.log',
  408. ])
  409. expect(fake.writtenFileData.get('/workspace/.dsh-e2b/processes/one/environment')).toBe(
  410. 'PATH=/bin\0UNICODE=你好\0HOME=/home/user\0FOO-BAR=hyphen-value\0--split-string=literal-value\0DEEPSEEK_API_KEY=explicit-secret\0DSH_MODE=test\0',
  411. )
  412. let piped = ''
  413. handle.stdout!.on('data', (chunk) => { piped += String(chunk) })
  414. await fake.stdout('pipe-data')
  415. await fake.stderr('err')
  416. fake.finish()
  417. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  418. expect(piped).toBe('pipe-data')
  419. expect(handle.collected.stderr!.readFrom(0)).toMatchObject({ text: 'err', lossy: false })
  420. expect(fake.removed).toContain('/workspace/.dsh-e2b/processes/one/stderr.log')
  421. await expect(handle.waitForExit()).resolves.toBe(true)
  422. })
  423. it('rejects an unrepresentable graceMs before any remote work', () => {
  424. const ctx = new Context()
  425. const service = Object.create(E2BSubprocessRuntime.prototype) as E2BSubprocessRuntime
  426. Reflect.set(service, 'disposing', false)
  427. Reflect.set(service, 'ctx', ctx)
  428. for (const graceMs of [0, -1, Number.NaN, Number.POSITIVE_INFINITY]) {
  429. expect(() => service.spawn(spec({ graceMs }))).toThrow('graceMs must be a positive finite number')
  430. void expect(service.spawnTerminal({
  431. argv: ['bash'], cwd: '/w', rows: 24, cols: 80, graceMs,
  432. })).rejects.toThrow('graceMs must be a positive finite number')
  433. }
  434. })
  435. it('rejects malformed environment entries before command start', async () => {
  436. for (const env of [{ 'BAD=NAME': 'x' }, { BAD: 'x\0INJECTED=1' }]) {
  437. const fake = new FakeSandbox()
  438. const handle = testHandle(runtime(fake), spec({ env }), '/runtime/invalid-environment')
  439. await expect(handle.done).rejects.toThrow('environment entries')
  440. expect(fake.startOptions).toBeUndefined()
  441. expect(fake.removed).toContain('/runtime/invalid-environment')
  442. }
  443. })
  444. it('preserves UTF-8 bytes when the ASCII transport is split across callbacks', async () => {
  445. const fake = new FakeSandbox()
  446. const handle = testHandle(runtime(fake), spec({
  447. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  448. }), '/runtime/split-utf8')
  449. await flush()
  450. const chunks: Buffer[] = []
  451. handle.stdout!.on('data', (chunk: Buffer) => { chunks.push(chunk) })
  452. for (const character of `${Buffer.from('A你好B').toString('base64')}\n`) {
  453. await fake.stdoutWire(character)
  454. }
  455. fake.finish()
  456. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  457. expect(Buffer.concat(chunks).toString('utf8')).toBe('A你好B')
  458. })
  459. it('rejects malformed output transport without confusing it with a consumer sink failure', async () => {
  460. const fake = new FakeSandbox()
  461. const handle = testHandle(runtime(fake), spec(), '/runtime/malformed-output')
  462. await flush()
  463. await fake.stdoutWire('%\n')
  464. fake.finish()
  465. await expect(handle.done).rejects.toThrow('invalid base64 output transport')
  466. const stderrFake = new FakeSandbox()
  467. const stderrHandle = testHandle(runtime(stderrFake), spec(), '/runtime/malformed-stderr')
  468. await flush()
  469. await stderrFake.stderrWire('%\n')
  470. stderrFake.finish()
  471. await expect(stderrHandle.done).rejects.toThrow('invalid base64 output transport')
  472. })
  473. it('rejects a naturally completed command whose encoder omits its completion frame', async () => {
  474. const fake = new FakeSandbox()
  475. const handle = testHandle(runtime(fake), spec(), '/runtime/incomplete-output')
  476. await flush()
  477. fake.alive = false
  478. fake.handle.succeed(0)
  479. await expect(handle.done).rejects.toThrow('incomplete output transport')
  480. })
  481. it('bounds descendant-held output draining and withholds the incomplete spill', async () => {
  482. const fake = new FakeSandbox()
  483. const handle = testHandle(runtime(fake), spec({ graceMs: 5 }), '/runtime/drain-bound')
  484. await flush()
  485. await fake.stdout('leader-output')
  486. fake.exitStatus = '0\n'
  487. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  488. expect(fake.handle.disconnects).toBe(1)
  489. expect(handle.collected.stdout?.readFrom(0)).toEqual({
  490. text: 'tput',
  491. nextOffset: 13,
  492. lossy: true,
  493. })
  494. expect(fake.removed).toContain('/runtime/drain-bound/stdout.log')
  495. handle.terminate()
  496. await expect(handle.waitForExit()).resolves.toBe(true)
  497. })
  498. it('releases an inherited-output callback blocked on host backpressure at drain expiry', async () => {
  499. const fake = new FakeSandbox()
  500. const written: string[] = []
  501. const stdoutWrite = vi.spyOn(process.stdout, 'write').mockImplementation(((chunk: Uint8Array) => {
  502. written.push(Buffer.from(chunk).toString())
  503. return false
  504. }) as typeof process.stdout.write)
  505. try {
  506. const handle = testHandle(runtime(fake), spec({
  507. graceMs: 5,
  508. stdio: { stdin: 'ignore', stdout: 'inherit', stderr: { maxBytes: 4 } },
  509. }), '/runtime/inherit-backpressure')
  510. await flush()
  511. let callbackSettled = false
  512. const blocked = fake.stdout('blocked bytes').then(() => { callbackSettled = true })
  513. await flush()
  514. expect(callbackSettled).toBe(false)
  515. fake.exitStatus = '0\n'
  516. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  517. await blocked
  518. expect(callbackSettled).toBe(true)
  519. expect(written.join('')).toBe('blocked bytes')
  520. expect(fake.handle.disconnects).toBe(1)
  521. handle.terminate()
  522. await expect(handle.waitForExit()).resolves.toBe(true)
  523. } finally {
  524. stdoutWrite.mockRestore()
  525. }
  526. })
  527. it('waits for lossless raw-pipe output after the direct status is published', async () => {
  528. const fake = new FakeSandbox()
  529. const handle = testHandle(runtime(fake), spec({
  530. graceMs: 1,
  531. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  532. }), '/runtime/pipe-drain')
  533. let output = ''
  534. handle.stdout!.on('data', (chunk) => { output += String(chunk) })
  535. await flush()
  536. fake.exitStatus = '0\n'
  537. let settled = false
  538. void handle.done.then(() => { settled = true })
  539. await new Promise(resolve => setTimeout(resolve, 50))
  540. expect(settled).toBe(false)
  541. expect(fake.handle.disconnects).toBe(0)
  542. expect(fake.statusReads).toBe(0)
  543. await fake.stdout('complete protocol frame')
  544. fake.finish()
  545. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  546. expect(output).toBe('complete protocol frame')
  547. expect(fake.statusReads).toBe(1)
  548. })
  549. it('accepts clean encoder completion inside the output-drain grace', async () => {
  550. const fake = new FakeSandbox()
  551. const handle = testHandle(runtime(fake), spec({ graceMs: 100 }), '/runtime/drain-complete')
  552. await flush()
  553. fake.exitStatus = '0\n'
  554. fake.afterStatusRead = () => {
  555. fake.afterStatusRead = undefined
  556. setTimeout(() => { fake.finish() }, 0)
  557. }
  558. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  559. expect(fake.handle.disconnects).toBe(0)
  560. })
  561. it('preserves a published exit code when requested termination outlives output draining', async () => {
  562. const fake = new FakeSandbox()
  563. fake.trapsTerm = true
  564. fake.delaysKill = true
  565. fake.delaysKillCompletion = true
  566. fake.sdkKillStops = false
  567. const handle = testHandle(runtime(fake), spec({ graceMs: 5 }), '/runtime/drain-signal')
  568. await flush()
  569. handle.terminate()
  570. await vi.waitFor(() => { expect(fake.commandsSeen).toContain('kill -KILL -- -4242') })
  571. fake.exitStatus = '0\n'
  572. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  573. expect(fake.handle.disconnects).toBe(1)
  574. fake.alive = false
  575. handle.terminate()
  576. await expect(handle.waitForExit()).resolves.toBe(true)
  577. })
  578. it('preserves a published nonzero exit code when termination settles the SDK inside the drain grace', async () => {
  579. const fake = new FakeSandbox()
  580. const handle = testHandle(runtime(fake), spec({ graceMs: 100 }), '/runtime/drain-signal-settled')
  581. await flush()
  582. fake.exitStatus = '7\n'
  583. fake.afterStatusRead = () => {
  584. fake.afterStatusRead = undefined
  585. handle.terminate()
  586. }
  587. await expect(handle.done).resolves.toEqual({ exitCode: 7, signal: null })
  588. await expect(handle.waitForExit()).resolves.toBe(true)
  589. })
  590. it('rejects an invalid direct-command exit status', async () => {
  591. const fake = new FakeSandbox()
  592. const handle = testHandle(runtime(fake), spec(), '/runtime/invalid-status')
  593. await flush()
  594. fake.exitStatus = '999\n'
  595. await expect(handle.done).rejects.toThrow('invalid exit code')
  596. handle.terminate()
  597. await expect(handle.waitForExit()).resolves.toBe(true)
  598. })
  599. it('rolls back a published process group before rejecting a monitoring failure', async () => {
  600. const fake = new FakeSandbox()
  601. fake.statusError = new Error('status transport failed')
  602. const handle = testHandle(runtime(fake), spec(), '/runtime/status-failure')
  603. await expect(handle.done).rejects.toThrow('status transport failed')
  604. expect(fake.commandsSeen).toContain('kill -TERM -- -4242')
  605. expect(fake.alive).toBe(false)
  606. await expect(handle.waitForExit()).resolves.toBe(true)
  607. const failed = new FakeSandbox()
  608. failed.statusError = new Error('status transport failed')
  609. failed.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  610. failed.handle.killError = new Error('SDK kill failed')
  611. const retained = testHandle(runtime(failed), spec({ graceMs: 1 }), '/runtime/status-cleanup-failure')
  612. await expect(retained.done).rejects.toThrow(
  613. 'command monitoring failed and process-group rollback did not reach quiescence',
  614. )
  615. expect(failed.alive).toBe(true)
  616. failed.handle.killError = undefined
  617. retained.terminate()
  618. await expect(retained.waitForExit()).resolves.toBe(true)
  619. // A state-cleanup failure on top preserves the rollback failure instead of
  620. // re-aggregating only the original monitoring error.
  621. const triple = new FakeSandbox()
  622. triple.statusError = new Error('status transport failed')
  623. triple.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  624. triple.handle.killError = new Error('SDK kill failed')
  625. triple.nextRemoveError = new Error('state cleanup failed')
  626. const tripleHandle = testHandle(runtime(triple), spec({ graceMs: 1 }), '/runtime/triple-failure')
  627. const failure = await tripleHandle.done.catch((error: unknown) => error as AggregateError)
  628. expect(failure).toBeInstanceOf(AggregateError)
  629. expect((failure as AggregateError).message).toContain('private state cleanup failed')
  630. const nested = (failure as AggregateError).errors[0] as AggregateError
  631. expect(nested.message).toContain('rollback did not reach quiescence')
  632. triple.handle.killError = undefined
  633. tripleHandle.terminate()
  634. await expect(tripleHandle.waitForExit()).resolves.toBe(true)
  635. })
  636. it('surfaces deferred piped-stdin write and close failures as stream errors', async () => {
  637. const writeFake = new FakeSandbox()
  638. writeFake.deferStart()
  639. vi.spyOn(writeFake.handle, 'sendStdin').mockRejectedValueOnce('stdin rejected')
  640. const writeHandle = testHandle(runtime(writeFake), spec({
  641. stdio: { stdin: 'pipe', stdout: { maxBytes: 4 }, stderr: { maxBytes: 4 } },
  642. }), '/runtime/stdin-write-error')
  643. const writeError = once(writeHandle.stdin!, 'error')
  644. writeHandle.stdin!.write('input')
  645. writeFake.releaseStart()
  646. await expect(writeError).resolves.toMatchObject([{ message: 'stdin rejected' }])
  647. writeFake.finish()
  648. await writeHandle.done
  649. const closeFake = new FakeSandbox()
  650. vi.spyOn(closeFake.handle, 'closeStdin').mockRejectedValueOnce(new Error('close rejected'))
  651. const closeHandle = testHandle(runtime(closeFake), spec({
  652. stdio: { stdin: 'pipe', stdout: { maxBytes: 4 }, stderr: { maxBytes: 4 } },
  653. }), '/runtime/stdin-close-error')
  654. await flush()
  655. const closeError = once(closeHandle.stdin!, 'error')
  656. closeHandle.stdin!.end()
  657. await expect(closeError).resolves.toMatchObject([{ message: 'close rejected' }])
  658. closeFake.finish()
  659. await closeHandle.done
  660. })
  661. it('collects bounded tails, retains valid spills, and maps natural nonzero exits', async () => {
  662. const fake = new FakeSandbox()
  663. const handle = testHandle(runtime(fake), spec({
  664. stdio: {
  665. stdin: { data: 'batch' },
  666. stdout: { maxBytes: 4, spill: { maxBytes: 16 } },
  667. stderr: { maxBytes: 3 },
  668. },
  669. }), '/runtime/two')
  670. await flush()
  671. await fake.stdout('abcdef')
  672. await fake.stderr('12345')
  673. fake.finish(7)
  674. await expect(handle.done).resolves.toEqual({ exitCode: 7, signal: null })
  675. expect(fake.handle.sent).toEqual(['batch'])
  676. expect(fake.handle.closes).toBe(1)
  677. expect(handle.collected.stdout!.readFrom(0)).toEqual({
  678. text: 'cdef',
  679. nextOffset: 6,
  680. lossy: true,
  681. spillPath: '/runtime/two/stdout.log',
  682. })
  683. expect(handle.collected.stderr!.readFrom(0)).toEqual({ text: '345', nextOffset: 5, lossy: true })
  684. expect(fake.removed).not.toContain('/runtime/two/stdout.log')
  685. })
  686. it('removes a spill once the complete stream exceeds its cap', async () => {
  687. const fake = new FakeSandbox()
  688. const handle = testHandle(runtime(fake), spec({
  689. stdio: { stdin: 'ignore', stdout: { maxBytes: 2, spill: { maxBytes: 3 } }, stderr: 'inherit' },
  690. }), '/runtime/oversize')
  691. await flush()
  692. await fake.stdout('abcd')
  693. await fake.stderr('')
  694. fake.finish()
  695. await handle.done
  696. expect(handle.collected.stdout!.readFrom(0)).toEqual({ text: 'cd', nextOffset: 4, lossy: true })
  697. expect(fake.removed).toContain('/runtime/oversize/stdout.log')
  698. const command = fake.commandsSeen.find(value => value.includes('dsh_e2b_tee='))!
  699. expect(command).toContain('"$dsh_e2b_head" -c 3')
  700. expect(command).toContain('/runtime/oversize/stdout.log')
  701. expect(command).toContain('"$dsh_e2b_tee" --output-error=warn-nopipe')
  702. expect(command).not.toContain('tee -a')
  703. })
  704. it('contains remote spill-removal failures and routes empty inherited output', async () => {
  705. const fake = new FakeSandbox()
  706. fake.nextRemoveError = new Error('already removed')
  707. const handle = testHandle(runtime(fake), spec({
  708. stdio: { stdin: 'ignore', stdout: 'inherit', stderr: { maxBytes: 4, spill: { maxBytes: 8 } } },
  709. }), '/runtime/remove-error')
  710. await flush()
  711. await fake.stdout('')
  712. await fake.stderr('')
  713. fake.finish()
  714. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  715. expect(fake.removed).toContain('/runtime/remove-error/stderr.log')
  716. })
  717. it('terminates a process group with TERM and reports the signal outcome', async () => {
  718. const fake = new FakeSandbox()
  719. const handle = testHandle(runtime(fake), spec(), '/runtime/term')
  720. await flush()
  721. handle.terminate()
  722. handle.terminate()
  723. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  724. await expect(handle.waitForExit()).resolves.toBe(true)
  725. expect(fake.commandsSeen).toContain('kill -TERM -- -4242')
  726. expect(fake.commandsSeen).not.toContain('kill -KILL -- -4242')
  727. const signals = fake.commandsSeen.filter(command => command.startsWith('kill -')).length
  728. fake.alive = true
  729. handle.terminate()
  730. await flush()
  731. expect(fake.alive).toBe(true)
  732. expect(fake.commandsSeen.filter(command => command.startsWith('kill -'))).toHaveLength(signals)
  733. })
  734. it('makes termination a permanent no-op after natural quiescence is observed', async () => {
  735. const fake = new FakeSandbox()
  736. const handle = testHandle(runtime(fake), spec(), '/runtime/natural-quiescence')
  737. await flush()
  738. fake.finish()
  739. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  740. await expect(handle.waitForExit()).resolves.toBe(true)
  741. const signals = fake.commandsSeen.filter(command => command.startsWith('kill -')).length
  742. fake.alive = true
  743. handle.terminate()
  744. await flush()
  745. expect(fake.alive).toBe(true)
  746. expect(fake.commandsSeen.filter(command => command.startsWith('kill -'))).toHaveLength(signals)
  747. })
  748. it('treats a zombie-only process group as quiescent', async () => {
  749. const fake = new FakeSandbox()
  750. fake.zombieOnly = true
  751. const handle = testHandle(runtime(fake), spec(), '/runtime/zombie-quiescence')
  752. await flush()
  753. await expect(handle.waitForExit()).resolves.toBe(true)
  754. expect(fake.commandsSeen).toContain(
  755. 'set -o pipefail; ps -eo pgid=,stat= | awk \'$1 == 4242 && $2 !~ /^[ZXx]/ { live=1 } END { if (live) print "live" }\'',
  756. )
  757. fake.finish()
  758. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  759. })
  760. it('keeps proven quiescence after a concurrent termination transport fails', async () => {
  761. const fake = new FakeSandbox()
  762. fake.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  763. fake.handle.killError = new Error('SDK kill failed')
  764. fake.deferSignals()
  765. const handle = testHandle(runtime(fake), spec(), '/runtime/quiescent-race')
  766. await flush()
  767. handle.terminate()
  768. await vi.waitFor(() => { expect(fake.commandsSeen).toContain('kill -TERM -- -4242') })
  769. fake.alive = false
  770. await expect(handle.waitForExit()).resolves.toBe(true)
  771. fake.probeError = new Error('post-quiescence probe failed')
  772. fake.releaseSignals()
  773. await vi.waitFor(() => { expect(fake.handle.kills).toBe(1) })
  774. fake.finish()
  775. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  776. const signals = fake.commandsSeen.filter(command => command.startsWith('kill -')).length
  777. handle.terminate()
  778. await expect(handle.waitForExit()).resolves.toBe(true)
  779. expect(fake.commandsSeen.filter(command => command.startsWith('kill -'))).toHaveLength(signals)
  780. })
  781. it('escalates a TERM-trapping process group to KILL and uses the SDK kill as fallback', async () => {
  782. const fake = new FakeSandbox()
  783. fake.trapsTerm = true
  784. fake.handle.killError = new Error('already gone')
  785. const handle = testHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/kill')
  786. await flush()
  787. handle.terminate()
  788. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  789. await expect(handle.waitForExit()).resolves.toBe(true)
  790. expect(fake.commandsSeen).toContain('kill -KILL -- -4242')
  791. expect(fake.handle.kills).toBe(1)
  792. })
  793. it('keeps force cleanup retryable until quiescence is proven', async () => {
  794. const fake = new FakeSandbox()
  795. fake.trapsTerm = true
  796. fake.delaysKill = true
  797. fake.delaysKillCompletion = true
  798. fake.sdkKillStops = false
  799. const handle = testHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/termination-fence')
  800. await flush()
  801. handle.terminate()
  802. await vi.waitFor(() => { expect(fake.handle.kills).toBe(1) })
  803. await expect(handle.waitForExit()).rejects.toThrow('remained live after force termination')
  804. expect(fake.alive).toBe(true)
  805. fake.delaysKill = false
  806. fake.delaysKillCompletion = false
  807. fake.sdkKillStops = true
  808. handle.terminate()
  809. await expect(handle.waitForExit()).resolves.toBe(true)
  810. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  811. })
  812. it('honors termination requested before asynchronous startup finishes', async () => {
  813. const fake = new FakeSandbox()
  814. fake.deferStart()
  815. const handle = testHandle(runtime(fake), spec(), '/runtime/deferred-kill')
  816. handle.terminate()
  817. fake.releaseStart()
  818. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  819. })
  820. it('aborts a stalled preparation request before reporting startup quiescence', async () => {
  821. const fake = new FakeSandbox()
  822. let preparationSignal: AbortSignal | undefined
  823. fake.environmentRequest = async (signal) => {
  824. preparationSignal = signal
  825. await new Promise<never>((_resolve, reject) => {
  826. const rejectAbort = (): void => {
  827. const reason: unknown = signal?.reason
  828. reject(reason instanceof Error ? reason : new Error(String(reason)))
  829. }
  830. if (signal?.aborted === true) {
  831. rejectAbort()
  832. return
  833. }
  834. signal?.addEventListener('abort', rejectAbort, { once: true })
  835. })
  836. }
  837. const handle = testHandle(runtime(fake), spec(), '/runtime/stalled-preparation')
  838. await vi.waitFor(() => { expect(preparationSignal).toBeDefined() })
  839. handle.terminate()
  840. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  841. await expect(handle.waitForExit()).resolves.toBe(true)
  842. expect(preparationSignal?.aborted).toBe(true)
  843. expect(fake.startOptions).toBeUndefined()
  844. })
  845. it('kills through the provisional SDK handle before process-group publication', async () => {
  846. const fake = new FakeSandbox()
  847. fake.deferProcessGroupRead()
  848. fake.signalErrors.push(commandError(1), commandError(1))
  849. const handle = testHandle(runtime(fake), spec(), '/runtime/pre-publication-kill')
  850. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  851. handle.terminate()
  852. await vi.waitFor(() => { expect(fake.handle.kills).toBe(1) })
  853. expect(fake.alive).toBe(false)
  854. await expect(handle.waitForExit()).resolves.toBe(true)
  855. fake.releaseProcessGroupRead()
  856. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  857. })
  858. it('does not treat an unsuccessful SDK fallback as provisional group quiescence', async () => {
  859. const fake = new FakeSandbox()
  860. fake.deferProcessGroupRead()
  861. fake.trapsTerm = true
  862. fake.delaysKill = true
  863. fake.delaysKillCompletion = true
  864. fake.sdkKillStops = false
  865. fake.handle.killResult = false
  866. const handle = testHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/provisional-sdk-false')
  867. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  868. handle.terminate()
  869. await vi.waitFor(() => { expect(fake.handle.kills).toBe(1) })
  870. await expect(handle.waitForExit()).rejects.toThrow('remained live after force termination')
  871. fake.alive = false
  872. handle.terminate()
  873. await expect(handle.waitForExit()).resolves.toBe(true)
  874. fake.releaseProcessGroupRead()
  875. fake.finish()
  876. await handle.done
  877. })
  878. it('bounds a quiescence observer while provisional termination is awaiting the controller', async () => {
  879. const fake = new FakeSandbox()
  880. fake.deferProcessGroupRead()
  881. const reconnect = Promise.withResolvers<Sandbox>()
  882. let calls = 0
  883. const delayedRuntime = runtime(fake, async () => {
  884. calls += 1
  885. return calls === 1 ? fake.sandbox : await reconnect.promise
  886. })
  887. const handle = testHandle(delayedRuntime, spec(), '/runtime/pre-publication-observer')
  888. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  889. handle.terminate()
  890. const controller = new AbortController()
  891. const waiting = handle.waitForExit(controller.signal)
  892. await flush()
  893. controller.abort()
  894. await expect(waiting).resolves.toBe(false)
  895. reconnect.resolve(fake.sandbox)
  896. await expect(handle.waitForExit()).resolves.toBe(true)
  897. fake.releaseProcessGroupRead()
  898. await handle.done
  899. })
  900. it('proves a provisional group exit when the SDK kill fallback fails', async () => {
  901. const fake = new FakeSandbox()
  902. fake.deferProcessGroupRead()
  903. fake.trapsTerm = true
  904. fake.handle.killError = new Error('SDK kill unavailable')
  905. const handle = testHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/pre-publication-group-kill')
  906. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  907. handle.terminate()
  908. await expect(handle.waitForExit()).resolves.toBe(true)
  909. fake.releaseProcessGroupRead()
  910. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  911. })
  912. it('reports failed provisional group and SDK force transports', async () => {
  913. const fake = new FakeSandbox()
  914. fake.deferProcessGroupRead()
  915. fake.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  916. fake.handle.killError = new Error('SDK kill failed')
  917. const handle = testHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/pre-publication-failure')
  918. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  919. handle.terminate()
  920. await expect(handle.waitForExit()).rejects.toThrow('remained live after force termination')
  921. fake.handle.killError = undefined
  922. handle.terminate()
  923. await expect(handle.waitForExit()).resolves.toBe(true)
  924. fake.releaseProcessGroupRead()
  925. await handle.done
  926. const absentGroup = new FakeSandbox()
  927. absentGroup.deferProcessGroupRead()
  928. absentGroup.signalErrors.push(commandError(1), commandError(1))
  929. absentGroup.handle.killError = new Error('SDK kill failed without a provisional group')
  930. const absentHandle = testHandle(
  931. runtime(absentGroup),
  932. spec({ graceMs: 1 }),
  933. '/runtime/pre-publication-absent-group',
  934. )
  935. await vi.waitFor(() => { expect(absentGroup.startOptions).toBeDefined() })
  936. absentHandle.terminate()
  937. await expect(absentHandle.waitForExit()).rejects.toThrow('remained live after force termination')
  938. absentGroup.handle.killError = undefined
  939. absentHandle.terminate()
  940. await expect(absentHandle.waitForExit()).resolves.toBe(true)
  941. absentGroup.releaseProcessGroupRead()
  942. await absentHandle.done
  943. const optimisticSdk = new FakeSandbox()
  944. optimisticSdk.deferProcessGroupRead()
  945. optimisticSdk.signalErrors.push(commandError(1), commandError(1))
  946. optimisticSdk.sdkKillStops = false
  947. const optimisticHandle = testHandle(
  948. runtime(optimisticSdk),
  949. spec({ graceMs: 1 }),
  950. '/runtime/pre-publication-optimistic-sdk',
  951. )
  952. await vi.waitFor(() => { expect(optimisticSdk.startOptions).toBeDefined() })
  953. optimisticHandle.terminate()
  954. await expect(optimisticHandle.waitForExit()).rejects.toThrow('remained live after force termination')
  955. optimisticHandle.terminate()
  956. await expect(optimisticHandle.waitForExit()).resolves.toBe(true)
  957. optimisticSdk.releaseProcessGroupRead()
  958. await optimisticHandle.done
  959. })
  960. it('honors an already-aborted signal when constructing the asynchronous handle directly', async () => {
  961. const fake = new FakeSandbox()
  962. const handle = testHandle(runtime(fake), spec({ signal: AbortSignal.abort('stop') }), '/runtime/pre-aborted')
  963. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  964. })
  965. it('reacts to a signal that aborts after the remote command has started', async () => {
  966. const fake = new FakeSandbox()
  967. const controller = new AbortController()
  968. const handle = testHandle(runtime(fake), spec({ signal: controller.signal }), '/runtime/live-abort')
  969. await flush()
  970. controller.abort('stop')
  971. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  972. })
  973. it('can terminate a surviving process group after the command leader settles', async () => {
  974. const fake = new FakeSandbox()
  975. const handle = testHandle(runtime(fake), spec(), '/runtime/surviving-group')
  976. await flush()
  977. await fake.completeOutput()
  978. fake.handle.succeed(0)
  979. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  980. expect(fake.alive).toBe(true)
  981. handle.terminate()
  982. await flush()
  983. const signaled = fake.commandsSeen.includes('kill -TERM -- -4242')
  984. if (!signaled) fake.finish()
  985. await expect(handle.waitForExit()).resolves.toBe(true)
  986. expect(signaled).toBe(true)
  987. })
  988. it('bounds waitForExit while startup or a live group is pending', async () => {
  989. const fake = new FakeSandbox()
  990. fake.deferStart()
  991. const handle = testHandle(runtime(fake), spec(), '/runtime/wait')
  992. const beforeStart = new AbortController()
  993. const pending = handle.waitForExit(beforeStart.signal)
  994. beforeStart.abort()
  995. await expect(pending).resolves.toBe(false)
  996. await expect(handle.waitForExit(AbortSignal.abort())).resolves.toBe(false)
  997. fake.releaseStart()
  998. await flush()
  999. const live = new AbortController()
  1000. const liveWait = handle.waitForExit(live.signal)
  1001. live.abort()
  1002. await expect(liveWait).resolves.toBe(false)
  1003. fake.finish()
  1004. await handle.done
  1005. const terminatingFake = new FakeSandbox()
  1006. terminatingFake.deferStart()
  1007. const terminating = testHandle(runtime(terminatingFake), spec(), '/runtime/wait-termination-start')
  1008. terminating.terminate()
  1009. const beforeHandle = new AbortController()
  1010. const handlePending = terminating.waitForExit(beforeHandle.signal)
  1011. beforeHandle.abort()
  1012. await expect(handlePending).resolves.toBe(false)
  1013. terminatingFake.releaseStart()
  1014. await terminating.done
  1015. })
  1016. it('bounds both sides of the liveness-poll abort race', async () => {
  1017. const fake = new FakeSandbox()
  1018. const handle = testHandle(runtime(fake), spec(), '/runtime/poll-abort')
  1019. await flush()
  1020. const beforeTick = new AbortController()
  1021. fake.afterProbe = () => { beforeTick.abort(); fake.afterProbe = undefined }
  1022. await expect(handle.waitForExit(beforeTick.signal)).resolves.toBe(false)
  1023. const duringTick = new AbortController()
  1024. fake.afterProbe = () => {
  1025. fake.afterProbe = undefined
  1026. setTimeout(() => { duringTick.abort() }, 0)
  1027. }
  1028. await expect(handle.waitForExit(duringTick.signal)).resolves.toBe(false)
  1029. const duringProbe = new AbortController()
  1030. fake.beforeProbe = () => { duringProbe.abort(); fake.beforeProbe = undefined }
  1031. await expect(handle.waitForExit(duringProbe.signal)).resolves.toBe(false)
  1032. let racedAbort = false
  1033. const raceSignal = {
  1034. get aborted() { return racedAbort },
  1035. addEventListener: () => { racedAbort = true },
  1036. removeEventListener: () => {},
  1037. } as unknown as AbortSignal
  1038. await expect(handle.waitForExit(raceSignal)).resolves.toBe(false)
  1039. fake.finish()
  1040. await handle.done
  1041. })
  1042. it('observes a live group across one successful bounded poll', async () => {
  1043. const fake = new FakeSandbox()
  1044. const handle = testHandle(runtime(fake), spec(), '/runtime/poll-success')
  1045. await flush()
  1046. setTimeout(() => { fake.finish() }, 1)
  1047. await expect(handle.waitForExit(new AbortController().signal)).resolves.toBe(true)
  1048. await handle.done
  1049. })
  1050. it('treats startup failure as no live tree and contains readiness rejection', async () => {
  1051. const fake = new FakeSandbox()
  1052. fake.backgroundError = new Error('start failed')
  1053. const handle = testHandle(runtime(fake), spec(), '/runtime/fail')
  1054. await expect(handle.done).rejects.toThrow('start failed')
  1055. expect(fake.removed).toContain('/runtime/fail/environment')
  1056. expect(fake.removed).toContain('/runtime/fail')
  1057. await expect(handle.waitForExit()).resolves.toBe(true)
  1058. handle.terminate()
  1059. const unavailableHandle = testHandle(
  1060. runtime(new FakeSandbox(), async () => { throw new Error('sandbox unavailable') }),
  1061. spec(),
  1062. '/runtime/unavailable-start',
  1063. )
  1064. await expect(unavailableHandle.done).rejects.toThrow('sandbox unavailable')
  1065. await expect(unavailableHandle.waitForExit()).resolves.toBe(true)
  1066. const envFailure = new FakeSandbox()
  1067. envFailure.envError = new Error('ambient lookup failed')
  1068. const envHandle = testHandle(runtime(envFailure), spec(), '/runtime/env-failure')
  1069. await expect(envHandle.done).rejects.toThrow('ambient lookup failed')
  1070. expect(envFailure.removed).toEqual([])
  1071. const expectEnvironmentFailure = async (name: string, wire: string, message: string): Promise<void> => {
  1072. const fake = new FakeSandbox()
  1073. fake.environmentWire = wire
  1074. const failed = testHandle(runtime(fake), spec(), `/runtime/${name}`)
  1075. await expect(failed.done).rejects.toThrow(message)
  1076. }
  1077. const encodedEnvironment = Buffer.from('PATH=/bin\0').toString('base64')
  1078. const encodedHome = Buffer.from('/home/user').toString('base64')
  1079. await expectEnvironmentFailure('malformed-frame', '%', 'invalid base64')
  1080. await expectEnvironmentFailure('malformed-base64', `${encodedHome}\n%`, 'invalid base64')
  1081. await expectEnvironmentFailure(
  1082. 'invalid-utf8-home',
  1083. `${Buffer.from([0xff]).toString('base64')}\n${encodedEnvironment}`,
  1084. 'not valid UTF-8',
  1085. )
  1086. await expectEnvironmentFailure(
  1087. 'invalid-utf8-environment',
  1088. `${encodedHome}\n${Buffer.from([0xff]).toString('base64')}`,
  1089. 'not valid UTF-8',
  1090. )
  1091. await expectEnvironmentFailure(
  1092. 'relative-home',
  1093. `${Buffer.from('home/user').toString('base64')}\n${encodedEnvironment}`,
  1094. 'remote login home is invalid',
  1095. )
  1096. await expectEnvironmentFailure(
  1097. 'nul-home',
  1098. `${Buffer.from('/home/user\0tail').toString('base64')}\n${encodedEnvironment}`,
  1099. 'remote login home is invalid',
  1100. )
  1101. const cleanupFailure = new FakeSandbox()
  1102. cleanupFailure.backgroundError = new Error('start failed before credential consumption')
  1103. cleanupFailure.nextRemoveError = new Error('credential cleanup failed')
  1104. const cleanupHandle = testHandle(runtime(cleanupFailure), spec(), '/runtime/cleanup-failure')
  1105. await expect(cleanupHandle.done).rejects.toThrow('command failed and private state cleanup failed')
  1106. const absentState = new FakeSandbox()
  1107. absentState.backgroundError = new Error('start failed after external cleanup')
  1108. absentState.nextRemoveError = new FileNotFoundError('already removed')
  1109. const absentHandle = testHandle(runtime(absentState), spec(), '/runtime/absent-state')
  1110. await expect(absentHandle.done).rejects.toThrow('start failed after external cleanup')
  1111. })
  1112. it('bounds a readiness rejection with a still-live caller signal', async () => {
  1113. const fake = new FakeSandbox()
  1114. fake.deferStart()
  1115. fake.backgroundError = new Error('start failed')
  1116. const handle = testHandle(runtime(fake), spec(), '/runtime/fail-with-signal')
  1117. const waiting = handle.waitForExit(new AbortController().signal)
  1118. fake.releaseStart()
  1119. await expect(handle.done).rejects.toThrow('start failed')
  1120. await expect(waiting).resolves.toBe(true)
  1121. })
  1122. it('propagates an unavailable sandbox unless the caller aborts the wait', async () => {
  1123. const fake = new FakeSandbox()
  1124. let calls = 0
  1125. const unavailable = runtime(fake, async () => {
  1126. calls += 1
  1127. if (calls === 1) return fake.sandbox
  1128. throw new Error('connection unavailable')
  1129. })
  1130. const handle = testHandle(unavailable, spec(), '/runtime/unavailable')
  1131. await flush()
  1132. await expect(handle.waitForExit()).rejects.toThrow('connection unavailable')
  1133. fake.finish()
  1134. await handle.done
  1135. })
  1136. it('returns false when the caller aborts while reconnecting for liveness', async () => {
  1137. const fake = new FakeSandbox()
  1138. const reconnect = Promise.withResolvers<Sandbox>()
  1139. let calls = 0
  1140. const unavailable = runtime(fake, async () => {
  1141. calls += 1
  1142. return calls === 1 ? fake.sandbox : await reconnect.promise
  1143. })
  1144. const handle = testHandle(unavailable, spec(), '/runtime/reconnect-abort')
  1145. await flush()
  1146. const controller = new AbortController()
  1147. const waiting = handle.waitForExit(controller.signal)
  1148. await flush()
  1149. controller.abort()
  1150. reconnect.reject(new Error('connection unavailable'))
  1151. await expect(waiting).resolves.toBe(false)
  1152. fake.finish()
  1153. await handle.done
  1154. })
  1155. it('returns false when a liveness request itself is aborted and surfaces other probe failures', async () => {
  1156. const fake = new FakeSandbox()
  1157. const handle = testHandle(runtime(fake), spec(), '/runtime/probe')
  1158. await flush()
  1159. const controller = new AbortController()
  1160. controller.abort()
  1161. await expect(handle.waitForExit(controller.signal)).resolves.toBe(false)
  1162. fake.probeError = new Error('probe failed')
  1163. await expect(handle.waitForExit()).rejects.toThrow('probe failed')
  1164. fake.finish()
  1165. await handle.done
  1166. })
  1167. it('treats a timeout-killed sandbox as quiescent during liveness probing', async () => {
  1168. const fake = new FakeSandbox()
  1169. const handle = testHandle(runtime(fake), spec(), '/runtime/expired-sandbox')
  1170. await flush()
  1171. fake.finish()
  1172. await handle.done
  1173. fake.probeError = new SandboxNotFoundError('sandbox expired')
  1174. await expect(handle.waitForExit()).resolves.toBe(true)
  1175. })
  1176. it('treats a missing sandbox handle as quiescent during liveness acquisition', async () => {
  1177. const fake = new FakeSandbox()
  1178. let calls = 0
  1179. const handle = testHandle(runtime(fake, async () => {
  1180. calls += 1
  1181. if (calls === 1) return fake.sandbox
  1182. throw new SandboxNotFoundError('sandbox expired')
  1183. }), spec(), '/runtime/expired-acquisition')
  1184. await flush()
  1185. await expect(handle.waitForExit()).resolves.toBe(true)
  1186. await fake.completeOutput()
  1187. fake.alive = false
  1188. fake.handle.succeed(0)
  1189. await handle.done
  1190. })
  1191. it('treats sandbox loss during termination as quiescent', async () => {
  1192. const fake = new FakeSandbox()
  1193. let calls = 0
  1194. const handle = testHandle(runtime(fake, async () => {
  1195. calls += 1
  1196. if (calls === 1) return fake.sandbox
  1197. throw new SandboxNotFoundError('sandbox expired')
  1198. }), spec(), '/runtime/expired-termination')
  1199. await flush()
  1200. await fake.completeOutput()
  1201. handle.terminate()
  1202. await expect(handle.waitForExit()).resolves.toBe(true)
  1203. fake.alive = false
  1204. fake.handle.succeed(0)
  1205. await handle.done
  1206. })
  1207. it('makes batch stdin close failures best-effort', async () => {
  1208. const fake = new FakeSandbox()
  1209. vi.spyOn(fake.handle, 'sendStdin').mockRejectedValueOnce(new Error('closed'))
  1210. const handle = testHandle(runtime(fake), spec({
  1211. stdio: { stdin: { data: 'ignored' }, stdout: { maxBytes: 4 }, stderr: { maxBytes: 4 } },
  1212. }), '/runtime/stdin-closed')
  1213. await flush()
  1214. fake.finish()
  1215. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  1216. })
  1217. it('rejects malformed SDK process ids and non-command settlement failures', async () => {
  1218. const invalidPid = new FakeSandbox()
  1219. invalidPid.handle.pid = 0
  1220. const invalid = testHandle(runtime(invalidPid), spec(), '/runtime/invalid-pid')
  1221. await expect(invalid.done).rejects.toThrow(/invalid command pid 0/)
  1222. expect(invalidPid.handle.kills).toBe(1)
  1223. expect(invalidPid.removed).toContain('/runtime/invalid-pid/environment')
  1224. await expect(invalid.waitForExit()).resolves.toBe(true)
  1225. const failedRollback = new FakeSandbox()
  1226. failedRollback.handle.pid = 0
  1227. failedRollback.handle.killError = new Error('invalid handle kill failed')
  1228. const retained = testHandle(runtime(failedRollback), spec(), '/runtime/invalid-pid-retained')
  1229. await expect(retained.done).rejects.toThrow('invalid command pid rollback did not reach quiescence')
  1230. await expect(retained.waitForExit()).rejects.toThrow('invalid handle kill failed')
  1231. failedRollback.handle.killError = undefined
  1232. retained.terminate()
  1233. await expect(retained.waitForExit()).resolves.toBe(true)
  1234. const crashedFake = new FakeSandbox()
  1235. const crashed = testHandle(runtime(crashedFake), spec(), '/runtime/crashed')
  1236. await flush()
  1237. crashedFake.alive = false
  1238. crashedFake.handle.crash(new Error('command transport failed'))
  1239. await expect(crashed.done).rejects.toThrow('command transport failed')
  1240. })
  1241. it('rejects invalid or absent process-group publication', async () => {
  1242. const invalidGroup = new FakeSandbox()
  1243. invalidGroup.processGroupId = 'not-a-pid\n'
  1244. invalidGroup.delaysKill = true
  1245. invalidGroup.sdkKillStops = false
  1246. invalidGroup.afterProbe = () => { invalidGroup.alive = false }
  1247. const invalid = testHandle(runtime(invalidGroup), spec(), '/runtime/invalid-group')
  1248. await expect(invalid.done).rejects.toThrow(/invalid process-group id/)
  1249. expect(invalidGroup.handle.kills).toBe(1)
  1250. expect(invalidGroup.commandsSeen).toContain('kill -KILL -- -4242')
  1251. await expect(invalid.waitForExit()).resolves.toBe(true)
  1252. // A rewritten pid file must not aim the kill at every process (`-- -1`).
  1253. const unsafeGroup = new FakeSandbox()
  1254. unsafeGroup.processGroupId = '1\n'
  1255. unsafeGroup.delaysKill = true
  1256. unsafeGroup.sdkKillStops = false
  1257. unsafeGroup.afterProbe = () => { unsafeGroup.alive = false }
  1258. const unsafe = testHandle(runtime(unsafeGroup), spec(), '/runtime/unsafe-group')
  1259. await expect(unsafe.done).rejects.toThrow(/unsafe published process-group id 1/)
  1260. expect(unsafeGroup.commandsSeen).not.toContain('kill -KILL -- -1')
  1261. await expect(unsafe.waitForExit()).resolves.toBe(true)
  1262. const absentGroup = new FakeSandbox()
  1263. absentGroup.processGroupId = ''
  1264. const absent = testHandle(runtime(absentGroup), spec(), '/runtime/absent-group')
  1265. await flush()
  1266. absentGroup.finish()
  1267. await expect(absent.done).rejects.toThrow(/exited before publishing/)
  1268. expect(absentGroup.handle.kills).toBe(1)
  1269. expect(absentGroup.commandsSeen).toContain('kill -KILL -- -4242')
  1270. await expect(absent.waitForExit()).resolves.toBe(true)
  1271. })
  1272. it('keeps polling while a running command has not published its process group yet', async () => {
  1273. const fake = new FakeSandbox()
  1274. fake.processGroupReads.push('', '4242\n')
  1275. const handle = testHandle(runtime(fake), spec(), '/runtime/delayed-group-publication', 1)
  1276. await vi.waitFor(() => { expect(fake.processGroupReads).toEqual([]) })
  1277. fake.finish()
  1278. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  1279. await expect(handle.waitForExit()).resolves.toBe(true)
  1280. })
  1281. it('preserves publication failure and reports cleanup that cannot be verified', async () => {
  1282. const fake = new FakeSandbox()
  1283. fake.processGroupId = 'not-a-pid\n'
  1284. fake.signalError = new Error('rollback signal failed')
  1285. fake.handle.killError = new Error('SDK kill failed')
  1286. const handle = testHandle(runtime(fake), spec(), '/runtime/failed-rollback')
  1287. let failure: unknown
  1288. try {
  1289. await handle.done
  1290. } catch (error: unknown) {
  1291. failure = error
  1292. }
  1293. expect(failure).toBeInstanceOf(AggregateError)
  1294. if (!(failure instanceof AggregateError)) throw new Error('expected AggregateError')
  1295. expect(failure.message).toBe('subprocess-e2b: process-group publication failed and rollback did not reach quiescence')
  1296. const failures = Array.from(failure.errors as Iterable<unknown>)
  1297. expect(failures).toHaveLength(2)
  1298. expect(failures[0]).toBeInstanceOf(Error)
  1299. expect(failures[1]).toBeInstanceOf(Error)
  1300. if (!(failures[0] instanceof Error) || !(failures[1] instanceof Error)) throw new Error('expected nested errors')
  1301. expect(failures[0].message).toContain('invalid process-group id')
  1302. expect(failures[1].message).toContain('remained live after force termination')
  1303. expect(fake.handle.kills).toBe(1)
  1304. const bounded = new AbortController()
  1305. const waiting = handle.waitForExit(bounded.signal)
  1306. bounded.abort()
  1307. await expect(waiting).resolves.toBe(false)
  1308. handle.terminate()
  1309. await expect(handle.waitForExit()).resolves.toBe(true)
  1310. expect(fake.commandsSeen).toContain('kill -TERM -- -4242')
  1311. const naturallyGone = new FakeSandbox()
  1312. naturallyGone.processGroupId = 'not-a-pid\n'
  1313. naturallyGone.signalError = new Error('rollback signal failed')
  1314. naturallyGone.handle.killError = new Error('SDK kill failed')
  1315. const observed = testHandle(runtime(naturallyGone), spec(), '/runtime/failed-rollback-observed')
  1316. await expect(observed.done).rejects.toThrow('process-group publication failed')
  1317. naturallyGone.alive = false
  1318. await expect(observed.waitForExit()).resolves.toBe(true)
  1319. })
  1320. it('handles output backpressure and contains a stderr sink failure', async () => {
  1321. const fake = new FakeSandbox()
  1322. const handle = testHandle(runtime(fake), spec({
  1323. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: 'pipe' },
  1324. }), '/runtime/backpressure')
  1325. await flush()
  1326. handle.stdout!.on('error', () => {})
  1327. const stdoutWrite = vi.spyOn(handle.stdout!, 'write').mockReturnValueOnce(false)
  1328. const stdoutPending = fake.stdout('blocked')
  1329. queueMicrotask(() => { handle.stdout!.emit('drain') })
  1330. await stdoutPending
  1331. stdoutWrite.mockRestore()
  1332. handle.stderr!.on('error', () => {})
  1333. const stderrWrite = vi.spyOn(handle.stderr!, 'write').mockReturnValueOnce(false)
  1334. const stderrPending = fake.stderr('broken')
  1335. queueMicrotask(() => { handle.stderr!.emit('error', new Error('sink failed')) })
  1336. await stderrPending
  1337. stderrWrite.mockRestore()
  1338. fake.finish()
  1339. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  1340. })
  1341. it('settles output backpressure when the consumer closes the pipe', async () => {
  1342. const fake = new FakeSandbox()
  1343. const handle = testHandle(runtime(fake), spec({
  1344. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  1345. }), '/runtime/backpressure-close')
  1346. await flush()
  1347. const stdoutWrite = vi.spyOn(handle.stdout!, 'write').mockReturnValueOnce(false)
  1348. const pending = fake.stdout('discarded')
  1349. queueMicrotask(() => { handle.stdout!.destroy() })
  1350. await pending
  1351. stdoutWrite.mockRestore()
  1352. fake.finish()
  1353. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  1354. })
  1355. it('breaks output backpressure when termination owns the command', async () => {
  1356. const fake = new FakeSandbox()
  1357. const handle = testHandle(runtime(fake), spec({
  1358. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  1359. }), '/runtime/backpressure-termination')
  1360. await flush()
  1361. const stdoutWrite = vi.spyOn(handle.stdout!, 'write').mockReturnValueOnce(false)
  1362. let released = false
  1363. const pending = fake.stdout('blocked').then(() => { released = true })
  1364. await Promise.resolve()
  1365. handle.terminate()
  1366. await flush()
  1367. const releasedByTermination = released
  1368. if (!released) handle.stdout!.emit('drain')
  1369. await pending
  1370. stdoutWrite.mockRestore()
  1371. await handle.done
  1372. expect(releasedByTermination).toBe(true)
  1373. })
  1374. it('settles backpressure when a synchronous pipe write starts termination', async () => {
  1375. const fake = new FakeSandbox()
  1376. const handle = testHandle(runtime(fake), spec({
  1377. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  1378. }), '/runtime/backpressure-synchronous-termination')
  1379. await flush()
  1380. const stdoutWrite = vi.spyOn(handle.stdout!, 'write').mockImplementationOnce(() => {
  1381. handle.terminate()
  1382. return false
  1383. })
  1384. await expect(fake.stdout('blocked')).resolves.toBeUndefined()
  1385. stdoutWrite.mockRestore()
  1386. await handle.done
  1387. })
  1388. it('contains a pipe callback failure instead of rejecting command settlement', async () => {
  1389. const fake = new FakeSandbox()
  1390. const handle = testHandle(runtime(fake), spec({
  1391. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: { maxBytes: 4 } },
  1392. }), '/runtime/pipe-error')
  1393. await flush()
  1394. const emitted = once(handle.stdout!, 'error')
  1395. handle.stdout!.destroy(new Error('consumer failed'))
  1396. await emitted
  1397. await fake.stdout('late output')
  1398. fake.finish()
  1399. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  1400. })
  1401. it('contains an already-gone group signal and escalates after a TERM transport failure', async () => {
  1402. const gone = new FakeSandbox()
  1403. gone.trapsTerm = true
  1404. gone.signalError = commandError(1)
  1405. const goneHandle = testHandle(runtime(gone), spec({ graceMs: 1 }), '/runtime/gone-signal')
  1406. await flush()
  1407. goneHandle.terminate()
  1408. await expect(goneHandle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  1409. const failed = new FakeSandbox()
  1410. failed.signalError = new Error('signal transport failed')
  1411. const failedHandle = testHandle(runtime(failed), spec(), '/runtime/failed-signal')
  1412. await flush()
  1413. failedHandle.terminate()
  1414. await expect(failedHandle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  1415. expect(failed.commandsSeen).toContain('kill -KILL -- -4242')
  1416. })
  1417. it('allows termination retry after both force transports fail', async () => {
  1418. const fake = new FakeSandbox()
  1419. fake.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  1420. fake.handle.killError = new Error('SDK kill failed')
  1421. const handle = testHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/retry-signal')
  1422. await flush()
  1423. handle.terminate()
  1424. await expect(handle.waitForExit()).rejects.toThrow('remained live after force termination')
  1425. fake.handle.killError = undefined
  1426. handle.terminate()
  1427. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  1428. expect(fake.commandsSeen.filter(command => command.startsWith('kill -TERM '))).toHaveLength(2)
  1429. const missingGroup = new FakeSandbox()
  1430. missingGroup.trapsTerm = true
  1431. missingGroup.signalErrors.push(undefined, commandError(1))
  1432. missingGroup.handle.killError = new Error('SDK kill failed after group exit race')
  1433. const raced = testHandle(runtime(missingGroup), spec({ graceMs: 1 }), '/runtime/group-exit-race')
  1434. await flush()
  1435. raced.terminate()
  1436. await expect(raced.waitForExit()).rejects.toThrow('remained live after force termination')
  1437. missingGroup.handle.killError = undefined
  1438. raced.terminate()
  1439. await expect(raced.waitForExit()).resolves.toBe(true)
  1440. })
  1441. it('rejects an optimistic SDK kill while descendants survive a failed group KILL', async () => {
  1442. const fake = new FakeSandbox()
  1443. fake.trapsTerm = true
  1444. fake.sdkKillStops = false
  1445. fake.signalErrors.push(undefined, new Error('KILL transport failed'))
  1446. const handle = testHandle(runtime(fake), spec({ graceMs: 1 }), '/runtime/optimistic-sdk-kill')
  1447. await flush()
  1448. handle.terminate()
  1449. await expect(handle.waitForExit()).rejects.toThrow('remained live after force termination')
  1450. expect(fake.alive).toBe(true)
  1451. fake.sdkKillStops = true
  1452. handle.terminate()
  1453. await expect(handle.waitForExit()).resolves.toBe(true)
  1454. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  1455. })
  1456. })
  1457. describe('E2BSubprocessRuntime', () => {
  1458. async function service(
  1459. fake = new FakeSandbox(),
  1460. providedRuntime: E2BRuntime = runtime(fake),
  1461. ): Promise<{ ctx: Context; fiber: Awaited<ReturnType<Context['plugin']>> }> {
  1462. const ctx = new Context()
  1463. ctx.provide('e2b', providedRuntime)
  1464. const fiber = await ctx.plugin(E2BSubprocessRuntime)
  1465. return { ctx, fiber }
  1466. }
  1467. it('registers handles and disposal terminates and joins live remote groups regardless of sandbox policy', async () => {
  1468. const fake = new FakeSandbox()
  1469. fake.trapsTerm = true
  1470. const { ctx, fiber } = await service(fake)
  1471. const handle = ctx.subprocess.spawn(spec({ graceMs: 1 }))
  1472. await flush()
  1473. await fiber.dispose()
  1474. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGKILL' })
  1475. expect(fake.alive).toBe(false)
  1476. })
  1477. it('awaits SDK settlement after the remote process group becomes quiescent', async () => {
  1478. const fake = new FakeSandbox()
  1479. fake.trapsTerm = true
  1480. const { ctx, fiber } = await service(fake)
  1481. const handle = ctx.subprocess.spawn(spec())
  1482. await flush()
  1483. fake.alive = false
  1484. let disposed = false
  1485. const disposing = fiber.dispose().then(() => { disposed = true })
  1486. await flush()
  1487. expect(disposed).toBe(false)
  1488. fake.finish()
  1489. await disposing
  1490. await expect(handle.done).resolves.toEqual({ exitCode: 0, signal: null })
  1491. })
  1492. it('reports a failed termination transaction from disposal instead of waiting on done', async () => {
  1493. const fake = new FakeSandbox()
  1494. fake.signalErrors.push(new Error('TERM transport failed'), new Error('KILL transport failed'))
  1495. fake.handle.killError = new Error('SDK kill failed')
  1496. const { ctx, fiber } = await service(fake)
  1497. const handle = ctx.subprocess.spawn(spec({ graceMs: 1 }))
  1498. await flush()
  1499. await expect(fiber.dispose()).resolves.toBeUndefined()
  1500. await expect(handle.waitForExit()).rejects.toThrow('remained live after force termination')
  1501. fake.handle.killError = undefined
  1502. handle.terminate()
  1503. await expect(handle.waitForExit()).resolves.toBe(true)
  1504. await expect(handle.done).resolves.toEqual({ exitCode: null, signal: 'SIGTERM' })
  1505. })
  1506. it('aggregates sibling cleanup failures instead of reporting only the first', async () => {
  1507. const { ctx, fiber } = await service()
  1508. const disposalErrors: unknown[] = []
  1509. ctx.logger.error = ((error: unknown) => { disposalErrors.push(error) }) as typeof ctx.logger.error
  1510. const first = {
  1511. terminate: vi.fn(),
  1512. waitForExit: vi.fn(async () => { throw new Error('first cleanup failed') }),
  1513. done: Promise.resolve({ exitCode: 0, signal: null }),
  1514. } as unknown as E2BSubprocessHandle
  1515. const second = {
  1516. terminate: vi.fn(),
  1517. waitForExit: vi.fn(async () => { throw new Error('second cleanup failed') }),
  1518. done: Promise.resolve({ exitCode: 0, signal: null }),
  1519. } as unknown as E2BSubprocessHandle
  1520. const live = (ctx.subprocess as unknown as { live: Set<E2BSubprocessHandle> }).live
  1521. live.add(first)
  1522. live.add(second)
  1523. await fiber.dispose()
  1524. const failure = disposalErrors[0]
  1525. expect(failure).toBeInstanceOf(AggregateError)
  1526. if (!(failure instanceof AggregateError)) throw new Error('expected AggregateError')
  1527. expect(failure.errors.map(error => (error as Error).message).sort()).toEqual([
  1528. 'first cleanup failed',
  1529. 'second cleanup failed',
  1530. ])
  1531. })
  1532. it('waits for every owned cleanup before reporting a disposal failure', async () => {
  1533. const { ctx, fiber } = await service()
  1534. const failed = {
  1535. terminate: vi.fn(),
  1536. waitForExit: vi.fn(async () => { throw new Error('cleanup failed') }),
  1537. done: Promise.resolve({ exitCode: 0, signal: null }),
  1538. } as unknown as E2BSubprocessHandle
  1539. let finishCleanup!: () => void
  1540. const cleanup = new Promise<boolean>((resolve) => {
  1541. finishCleanup = () => { resolve(true) }
  1542. })
  1543. const draining = {
  1544. terminate: vi.fn(),
  1545. waitForExit: vi.fn(() => cleanup),
  1546. done: Promise.resolve({ exitCode: 0, signal: null }),
  1547. } as unknown as E2BSubprocessHandle
  1548. const live = (ctx.subprocess as unknown as { live: Set<E2BSubprocessHandle> }).live
  1549. live.add(failed)
  1550. live.add(draining)
  1551. let disposed = false
  1552. const disposing = fiber.dispose().then(() => { disposed = true })
  1553. await flush()
  1554. expect(disposed).toBe(false)
  1555. finishCleanup()
  1556. await disposing
  1557. expect(live).toEqual(new Set([failed]))
  1558. })
  1559. it('releases naturally settled handles before later service disposal', async () => {
  1560. const fake = new FakeSandbox()
  1561. const { ctx, fiber } = await service(fake)
  1562. const handle = ctx.subprocess.spawn(spec())
  1563. await flush()
  1564. fake.finish()
  1565. await handle.done
  1566. await flush()
  1567. const signalsBefore = fake.commandsSeen.filter(command => command.startsWith('kill -')).length
  1568. await fiber.dispose()
  1569. expect(fake.commandsSeen.filter(command => command.startsWith('kill -')).length).toBe(signalsBefore)
  1570. })
  1571. it('contains a release liveness failure and retries quiescence during disposal', async () => {
  1572. const fake = new FakeSandbox()
  1573. let calls = 0
  1574. const reconnecting = runtime(fake, async () => {
  1575. calls += 1
  1576. if (calls === 2) throw new Error('transient liveness failure')
  1577. return fake.sandbox
  1578. })
  1579. const { ctx, fiber } = await service(fake, reconnecting)
  1580. const handle = ctx.subprocess.spawn(spec())
  1581. await flush()
  1582. fake.finish()
  1583. await handle.done
  1584. await flush()
  1585. await fiber.dispose()
  1586. expect(calls).toBeGreaterThanOrEqual(3)
  1587. })
  1588. it('contains spawn rejection while disposal is joining the pending handle', async () => {
  1589. const fake = new FakeSandbox()
  1590. fake.deferStart()
  1591. fake.backgroundError = new Error('start failed during disposal')
  1592. const { ctx, fiber } = await service(fake)
  1593. const subprocess = ctx.subprocess
  1594. const handle = subprocess.spawn(spec())
  1595. await vi.waitFor(() => { expect(fake.startOptions).toBeDefined() })
  1596. const disposing = fiber.dispose()
  1597. await flush()
  1598. expect(() => subprocess.spawn(spec())).toThrow('service is disposing')
  1599. fake.releaseStart()
  1600. await expect(disposing).resolves.toBeUndefined()
  1601. await expect(handle.done).rejects.toThrow('start failed during disposal')
  1602. })
  1603. it('validates synchronous spawn preconditions', async () => {
  1604. const { ctx } = await service()
  1605. expect(() => ctx.subprocess.spawn(spec({ argv: [] }))).toThrow(/non-empty program/)
  1606. expect(() => ctx.subprocess.spawn(spec({ signal: AbortSignal.abort('stop') }))).toThrow(/aborted before spawn/)
  1607. })
  1608. it('registers the package-owned empty invariant installer', async () => {
  1609. const ctx = new Context()
  1610. await ctx.plugin(InvariantRegistry, { enabled: true })
  1611. const fiber = await ctx.plugin(E2BSubprocessInvariant).await()
  1612. await fiber.dispose()
  1613. })
  1614. })