index.ts 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208
  1. /**
  2. * E2B Service Provider for the subprocess capability seam. Each handle starts through the
  3. * shared sandbox and retains command output/status paths in that remote world.
  4. * @module @deepseek-ai/dsh-subprocess-e2b
  5. */
  6. import { randomUUID } from 'node:crypto'
  7. import { posix } from 'node:path'
  8. import { Context } from '@deepseek-ai/cordis'
  9. import z from '@deepseek-ai/schemastery'
  10. import { SubprocessRuntime } from '@deepseek-ai/dsh-subprocess'
  11. import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
  12. import type {
  13. SubprocessHandle,
  14. SubprocessSpawnSpec,
  15. SubprocessTerminalHandle,
  16. SubprocessTerminalSpawnSpec,
  17. } from '@deepseek-ai/dsh-subprocess'
  18. import { e2bControlEnvs, quoteE2BShellArg } from '@deepseek-ai/dsh-e2b'
  19. import { E2BSubprocessHandle } from './process.ts'
  20. import { asError, signalOpts } from './remote.ts'
  21. import { spawnE2BTerminal } from './terminal.ts'
  22. /** Configuration for the E2B subprocess adapter. */
  23. export interface Config {
  24. /** Remote status/liveness poll cadence in milliseconds; each tick is one control-plane request. */
  25. pollMs?: number
  26. }
  27. interface SchemaResolvedConfig extends Config {
  28. pollMs: number
  29. }
  30. interface TerminalSetup {
  31. done: Promise<void>
  32. controller: AbortController
  33. }
  34. /**
  35. * Enforce the seam's documented grace bound (positive, finite, one Node timer),
  36. * matching subprocess-local's spawn-time check; an unbounded grace would make
  37. * the remote force-escalation deadline unreachable.
  38. * @param graceMs - The spec's cleanup grace in milliseconds.
  39. */
  40. function requireRepresentableGrace(graceMs: number): void {
  41. if (!Number.isFinite(graceMs) || graceMs <= 0 || graceMs > MAX_TIMER_DELAY_MS) {
  42. throw new Error(`subprocess graceMs must be a positive finite number no greater than ${MAX_TIMER_DELAY_MS}`)
  43. }
  44. }
  45. /** E2B command manager registered as `ctx.subprocess`. */
  46. export class E2BSubprocessRuntime extends SubprocessRuntime {
  47. static inject = ['e2b']
  48. static Config: z<Config> = z.object({
  49. pollMs: z.number().default(20),
  50. })
  51. private readonly live = new Set<E2BSubprocessHandle>()
  52. private readonly terminals = new Set<SubprocessTerminalHandle>()
  53. private readonly terminalSetups = new Set<TerminalSetup>()
  54. private readonly pollMs: number
  55. private disposing = false
  56. /** Create the E2B subprocess service and bind its disposal policy. */
  57. constructor(ctx: Context, config: Config) {
  58. super(ctx)
  59. // Schemastery fills pollMs before construction; the type does not encode that step.
  60. const { pollMs } = config as SchemaResolvedConfig
  61. if (!Number.isSafeInteger(pollMs) || pollMs <= 0) {
  62. throw new Error('subprocess-e2b: pollMs must be a positive safe integer')
  63. }
  64. this.pollMs = pollMs
  65. ctx.effect(() => async () => {
  66. this.disposing = true
  67. for (const setup of this.terminalSetups) {
  68. setup.controller.abort(new Error('subprocess-e2b: service disposed during terminal setup'))
  69. }
  70. await Promise.all([...this.terminalSetups].map(setup => setup.done))
  71. const handles = [...this.live]
  72. const terminals = [...this.terminals]
  73. const pending: Promise<unknown>[] = []
  74. for (const handle of handles) {
  75. handle.terminate()
  76. pending.push(handle.waitForExit().then(async () => {
  77. await handle.done.catch(() => undefined)
  78. this.live.delete(handle)
  79. }))
  80. }
  81. for (const terminal of terminals) {
  82. pending.push(terminal.terminate().then(() => { this.terminals.delete(terminal) }))
  83. }
  84. const outcomes = await Promise.allSettled(pending)
  85. const failures = outcomes.flatMap<unknown>(outcome => outcome.status === 'rejected'
  86. ? [outcome.reason as unknown]
  87. : [])
  88. if (failures.length === 1) throw asError(failures[0])
  89. if (failures.length > 1) throw new AggregateError(failures, 'subprocess-e2b: teardown failed')
  90. }, 'e2b subprocess teardown')
  91. }
  92. /** @inheritdoc */
  93. async resolveExecutable(
  94. command: string,
  95. env?: Readonly<Record<string, string>>,
  96. signal?: AbortSignal,
  97. ): Promise<string> {
  98. if (command.length === 0) throw new Error('subprocess-e2b: executable name must be non-empty')
  99. signal?.throwIfAborted()
  100. const sandbox = await this.ctx.e2b.getSandbox()
  101. if (posix.isAbsolute(command)) {
  102. await sandbox.commands.run(
  103. `test -f ${quoteE2BShellArg(command)} -a -x ${quoteE2BShellArg(command)}`,
  104. { envs: e2bControlEnvs(), ...signalOpts(signal) },
  105. )
  106. signal?.throwIfAborted()
  107. return command
  108. }
  109. if (command.includes('/')) {
  110. throw new Error(
  111. `subprocess-e2b: command ${JSON.stringify(command)} is a relative path; use an absolute path or a bare PATH name`,
  112. )
  113. }
  114. const path = env?.PATH
  115. const prefix = path === undefined ? '' : `PATH=${quoteE2BShellArg(path)} `
  116. const result = await sandbox.commands.run(
  117. `${prefix}command -v -- ${quoteE2BShellArg(command)}`,
  118. { cwd: this.ctx.e2b.cwd, envs: e2bControlEnvs(), ...signalOpts(signal) },
  119. )
  120. signal?.throwIfAborted()
  121. const executable = result.stdout.trim()
  122. if (executable.includes('\n') || (!posix.isAbsolute(executable) && !executable.includes('/'))) {
  123. throw new Error(`subprocess-e2b: executable ${JSON.stringify(command)} did not resolve to one absolute path`)
  124. }
  125. // A relative result comes from a relative PATH entry; the lookup ran with the shared cwd.
  126. return posix.resolve(this.ctx.e2b.cwd, executable)
  127. }
  128. /** @inheritdoc */
  129. spawn(spec: SubprocessSpawnSpec): SubprocessHandle {
  130. if (this.disposing) throw new Error('subprocess-e2b: service is disposing')
  131. const program = spec.argv[0]
  132. if (program === undefined || program.length === 0) {
  133. throw new Error('invalid argv: expected a non-empty program name at argv[0]')
  134. }
  135. requireRepresentableGrace(spec.graceMs)
  136. if (spec.signal?.aborted === true) {
  137. throw new Error(`aborted before spawn: ${String(spec.signal.reason)}`)
  138. }
  139. const stateDir = posix.join(this.ctx.e2b.runtimeRoot, 'processes', randomUUID())
  140. const handle = new E2BSubprocessHandle(this.ctx.e2b, spec, stateDir, this.pollMs)
  141. this.live.add(handle)
  142. const release = async (): Promise<void> => {
  143. await handle.waitForExit()
  144. this.live.delete(handle)
  145. }
  146. void handle.done.then(release, release).catch((_automaticReleaseFailure: unknown) => {
  147. // Retain the handle so service disposal can retry its cleanup transaction.
  148. })
  149. return handle
  150. }
  151. /** @inheritdoc */
  152. async spawnTerminal(spec: SubprocessTerminalSpawnSpec): Promise<SubprocessTerminalHandle> {
  153. if (this.disposing) throw new Error('subprocess-e2b: service is disposing')
  154. const program = spec.argv[0]
  155. if (program === undefined || program.length === 0) {
  156. throw new Error('subprocess-e2b: terminal argv must contain a program')
  157. }
  158. requireRepresentableGrace(spec.graceMs)
  159. spec.signal?.throwIfAborted()
  160. const stateDir = posix.join(this.ctx.e2b.runtimeRoot, 'terminals', randomUUID())
  161. const done = Promise.withResolvers<void>()
  162. const setup: TerminalSetup = { done: done.promise, controller: new AbortController() }
  163. const setupSignal = spec.signal === undefined
  164. ? setup.controller.signal
  165. : AbortSignal.any([spec.signal, setup.controller.signal])
  166. this.terminalSetups.add(setup)
  167. try {
  168. const terminal = await spawnE2BTerminal(
  169. this.ctx.e2b,
  170. { ...spec, signal: setupSignal },
  171. stateDir,
  172. this.pollMs,
  173. )
  174. this.terminals.add(terminal)
  175. // oxlint-disable-next-line typescript/no-unnecessary-condition -- Remote allocation yields to disposal.
  176. if (this.disposing) {
  177. await terminal.terminate()
  178. this.terminals.delete(terminal)
  179. throw new Error('subprocess-e2b: service disposed during terminal setup')
  180. }
  181. const release = async (): Promise<void> => {
  182. await terminal.terminate()
  183. this.terminals.delete(terminal)
  184. }
  185. void terminal.done.then(release, release).catch((_automaticReleaseFailure: unknown) => {
  186. // Retain the terminal so service disposal can retry its cleanup transaction.
  187. })
  188. return terminal
  189. } finally {
  190. this.terminalSetups.delete(setup)
  191. done.resolve()
  192. }
  193. }
  194. }
  195. export default E2BSubprocessRuntime