stream-wait-lifecycle.spec.ts 8.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190
  1. /** Caller-owned stream closure and bounded observations across remote allocation. */
  2. import { duplexPair, type Duplex } from 'node:stream'
  3. import { getEventListeners } from 'node:events'
  4. import { Context } from '@deepseek-ai/cordis'
  5. import { describe, expect, it, onTestFinished } from 'vitest'
  6. import { z } from 'zod'
  7. import type { SubprocessSpawnSpec } from '@deepseek-ai/dsh-subprocess'
  8. import { SshSubprocessRuntime } from '../src/index.ts'
  9. type Stage = 'prepare' | 'connect' | 'start' | 'wait' | 'terminate'
  10. const id = '8fbfb59c-fbce-44f7-bd87-5ab5a0ea02aa'
  11. const spec: SubprocessSpawnSpec = {
  12. argv: ['node'], cwd: '/workspace', graceMs: 5,
  13. stdio: { stdin: 'ignore', stdout: 'pipe', stderr: 'pipe' },
  14. }
  15. const finishedValue = { outcome: { exitCode: 0, signal: null }, spills: {}, collected: {} }
  16. async function setup(options: { pause?: Stage; failStart?: Error; failWait?: Error; failTerminate?: Error } = {}) {
  17. const ctx = new Context()
  18. const gate = Promise.withResolvers<undefined>()
  19. const entered = Promise.withResolvers<undefined>()
  20. const started = Promise.withResolvers<undefined>()
  21. const finished = Promise.withResolvers<typeof finishedValue>()
  22. const peers = new Map<string, { host: Duplex; remote: Duplex }>()
  23. const calls: string[] = []
  24. const errors: unknown[] = []
  25. ctx.logger.error = ((error: unknown) => { errors.push(error) }) as typeof ctx.logger.error
  26. const stage = async (name: Stage, signal?: AbortSignal) => {
  27. if (options.pause !== name) { signal?.throwIfAborted(); return }
  28. entered.resolve(undefined)
  29. signal?.throwIfAborted()
  30. const cancelled = Promise.withResolvers<never>()
  31. const abort = (): void => { cancelled.reject(new Error('SSH operation cancelled')) }
  32. signal?.addEventListener('abort', abort, { once: true })
  33. try { await Promise.race([gate.promise, cancelled.promise]) }
  34. finally { signal?.removeEventListener('abort', abort) }
  35. }
  36. const connection = {
  37. request: async <T>(method: string, _params: unknown, schema: z.ZodType<T>, signal?: AbortSignal): Promise<T> => {
  38. calls.push(method)
  39. let value: unknown
  40. if (method === 'process.prepare') {
  41. await stage('prepare', signal)
  42. value = { id, streams: {
  43. stdout: { path: '/tmp/stdout', capability: 'a'.repeat(64) },
  44. stderr: { path: '/tmp/stderr', capability: 'b'.repeat(64) },
  45. } }
  46. } else if (method === 'process.start') {
  47. await stage('start', signal)
  48. if (options.failStart !== undefined) throw options.failStart
  49. started.resolve(undefined)
  50. value = {}
  51. } else if (method === 'process.done') value = await finished.promise
  52. else if (method === 'process.wait') {
  53. await stage('wait', signal)
  54. if (options.failWait !== undefined) throw options.failWait
  55. value = true
  56. } else if (method === 'process.terminate') {
  57. await stage('terminate', signal)
  58. if (options.failTerminate !== undefined) throw options.failTerminate
  59. value = null
  60. } else throw new Error(`Unexpected request ${method}`)
  61. return schema.parse(value)
  62. },
  63. connectStream: async (endpoint: { path: string }, signal?: AbortSignal) => {
  64. await stage('connect', signal)
  65. const [host, remote] = duplexPair({ allowHalfOpen: true })
  66. host.on('error', () => {})
  67. remote.on('error', () => {})
  68. peers.set(endpoint.path, { host, remote })
  69. return host
  70. },
  71. dispose: () => Promise.resolve(),
  72. }
  73. ctx.provide('ssh', connection as never)
  74. const fiber = await ctx.plugin(SshSubprocessRuntime)
  75. onTestFinished(async () => {
  76. gate.resolve(undefined)
  77. finished.resolve(finishedValue)
  78. for (const peer of peers.values()) { peer.host.destroy(); peer.remote.destroy() }
  79. await fiber.dispose()
  80. if (options.failTerminate === undefined) expect(errors).toEqual([])
  81. })
  82. const peer = (name: string) => {
  83. const entry = peers.get(`/tmp/${name}`)
  84. if (entry === undefined) throw new Error('transport was not allocated')
  85. return entry
  86. }
  87. return { runtime: ctx.subprocess, calls, peer, entered: entered.promise, started: started.promise,
  88. release: () => { gate.resolve(undefined) }, finished }
  89. }
  90. describe('SSH public output streams', () => {
  91. it.each(['stdout', 'stderr'] as const)('closes only the matching transport when %s is destroyed', async (name) => {
  92. const test = await setup()
  93. const handle = test.runtime.spawn(spec)
  94. await test.started
  95. handle[name]!.destroy()
  96. await new Promise(resolve => setImmediate(resolve))
  97. expect(test.peer(name).host.destroyed).toBe(true)
  98. const other = name === 'stdout' ? 'stderr' : 'stdout'
  99. expect(test.peer(other).host.destroyed).toBe(false)
  100. expect(test.calls).not.toContain('process.terminate')
  101. })
  102. it.each(['stdout', 'stderr'] as const)('remembers %s closure before remote allocation finishes', async (name) => {
  103. const test = await setup({ pause: 'prepare' })
  104. const handle = test.runtime.spawn(spec)
  105. await test.entered
  106. handle[name]!.destroy()
  107. test.release()
  108. await test.started
  109. expect(test.peer(name).host.destroyed).toBe(true)
  110. expect(test.peer(name === 'stdout' ? 'stderr' : 'stdout').host.destroyed).toBe(false)
  111. })
  112. })
  113. describe('SSH range-observation cancellation', () => {
  114. it.each(['prepare', 'connect', 'start', 'wait', 'terminate'] as const)('bounds the whole observation while %s is pending', async (pause) => {
  115. const test = await setup({ pause })
  116. const handle = test.runtime.spawn(spec)
  117. if (pause === 'terminate') { await test.started; handle.terminate() }
  118. const controller = new AbortController()
  119. let settled = false
  120. const observed = handle.waitForExit(controller.signal).then(
  121. value => ({ value }), (error: unknown) => ({ error }),
  122. ).finally(() => { settled = true })
  123. await test.entered
  124. controller.abort('wait deadline')
  125. await new Promise(resolve => setImmediate(resolve))
  126. expect(settled).toBe(true)
  127. expect(await observed).toEqual({ value: false })
  128. expect(getEventListeners(controller.signal, 'abort')).toHaveLength(0)
  129. if (pause !== 'terminate') expect(test.calls).not.toContain('process.terminate')
  130. test.release()
  131. expect(await handle.waitForExit()).toBe(true)
  132. })
  133. it('bounds termination after failed startup while retaining real termination failure', async () => {
  134. const cleanup = new Error('range owner unavailable')
  135. const test = await setup({ pause: 'terminate', failStart: new Error('start refused'), failTerminate: cleanup })
  136. const handle = test.runtime.spawn(spec)
  137. const controller = new AbortController()
  138. let settled = false
  139. const observed = handle.waitForExit(controller.signal).then(
  140. value => ({ value }), (error: unknown) => ({ error }),
  141. ).finally(() => { settled = true })
  142. await test.entered
  143. controller.abort()
  144. await new Promise(resolve => setImmediate(resolve))
  145. expect(settled).toBe(true)
  146. expect(await observed).toEqual({ value: false })
  147. test.release()
  148. await expect(handle.waitForExit()).rejects.toBe(cleanup)
  149. })
  150. it('returns false for a pre-aborted live observation and true once quiescence is known', async () => {
  151. const test = await setup({ pause: 'prepare' })
  152. const handle = test.runtime.spawn(spec)
  153. const signal = AbortSignal.abort()
  154. let settled = false
  155. const observed = handle.waitForExit(signal).then(
  156. value => ({ value }), (error: unknown) => ({ error }),
  157. ).finally(() => { settled = true })
  158. await new Promise(resolve => setImmediate(resolve))
  159. expect(settled).toBe(true)
  160. expect(await observed).toEqual({ value: false })
  161. test.release()
  162. expect(await handle.waitForExit()).toBe(true)
  163. expect(await handle.waitForExit(signal)).toBe(true)
  164. })
  165. it('preserves real observation failure and detaches the caller cancellation listener', async () => {
  166. const failure = new Error('remote range observation failed')
  167. const test = await setup({ failWait: failure })
  168. const handle = test.runtime.spawn(spec)
  169. await test.started
  170. const controller = new AbortController()
  171. await expect(handle.waitForExit(controller.signal)).rejects.toBe(failure)
  172. expect(getEventListeners(controller.signal, 'abort')).toHaveLength(0)
  173. })
  174. it('returns true when an uncancelled bounded observation completes', async () => {
  175. const test = await setup()
  176. const handle = test.runtime.spawn(spec)
  177. const controller = new AbortController()
  178. expect(await handle.waitForExit(controller.signal)).toBe(true)
  179. expect(getEventListeners(controller.signal, 'abort')).toHaveLength(0)
  180. })
  181. })