control-retry.client.spec.ts 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347
  1. import { describe, expect, it, vi } from 'vitest'
  2. import type { ConnectionHandle } from '@deepseek-ai/dsh-api-remotes/client'
  3. import {
  4. RemoteStreamCarrierError,
  5. RemoteStream,
  6. } from '../src/client/index.ts'
  7. const DESCRIPTION = {
  8. version: 'fixture',
  9. cwd: '/fixture',
  10. attachedSessions: 0,
  11. home: '/home/fixture',
  12. canOpenPath: true,
  13. }
  14. function hostSource(initiallyAvailable: boolean): {
  15. connection: Pick<ConnectionHandle, 'hostDescription'>
  16. publish(available: boolean): void
  17. } {
  18. let current = initiallyAvailable ? DESCRIPTION : undefined
  19. const listeners = new Set<() => void>()
  20. return {
  21. connection: {
  22. hostDescription: {
  23. getSnapshot: () => current,
  24. subscribe: (listener) => {
  25. listeners.add(listener)
  26. return () => { listeners.delete(listener) }
  27. },
  28. },
  29. },
  30. publish: (available) => {
  31. current = available ? DESCRIPTION : undefined
  32. for (const listener of listeners) listener()
  33. },
  34. }
  35. }
  36. interface Generation<Item> {
  37. readonly values?: readonly (Item | Promise<Item>)[]
  38. readonly terminal?: Error
  39. readonly hold?: boolean
  40. readonly afterAbortError?: Error
  41. readonly close?: () => Promise<void>
  42. }
  43. function scripted<Item>(generations: Generation<Item>[], opened?: () => void) {
  44. return (signal: AbortSignal): AsyncIterable<Item> => ({
  45. async * [Symbol.asyncIterator](): AsyncIterator<Item> {
  46. const generation = generations.shift()
  47. if (generation === undefined) throw new Error('fixture has no stream generation')
  48. opened?.()
  49. try {
  50. for (const value of generation.values ?? []) yield await value
  51. if (generation.terminal !== undefined) throw generation.terminal
  52. if (generation.hold === true && !signal.aborted) {
  53. await new Promise<void>((resolve) => {
  54. signal.addEventListener('abort', () => { resolve() }, { once: true })
  55. })
  56. }
  57. if (generation.afterAbortError !== undefined) throw generation.afterAbortError
  58. } finally {
  59. await generation.close?.()
  60. }
  61. },
  62. })
  63. }
  64. function supervisor<Item>(
  65. connection: Pick<ConnectionHandle, 'hostDescription'>,
  66. generations: Generation<Item>[],
  67. carrierFailed?: (error: RemoteStreamCarrierError) => void,
  68. ): RemoteStream<Item> {
  69. return new RemoteStream(connection, {
  70. name: 'fixture stream',
  71. open: scripted(generations),
  72. ended: accepted => accepted
  73. ? new RemoteStreamCarrierError('accepted generation ended')
  74. : new Error('generation ended before acceptance'),
  75. ...(carrierFailed === undefined ? {} : { carrierFailed }),
  76. })
  77. }
  78. describe('RemoteStream', () => {
  79. it('annotates replacement generations and resets retry state after acceptance', async () => {
  80. const source = hostSource(true)
  81. const stream = supervisor(source.connection, [
  82. { values: ['first'], terminal: new RemoteStreamCarrierError('first lost') },
  83. { values: ['second'], hold: true },
  84. ])
  85. const iterator = stream[Symbol.asyncIterator]()
  86. const first = await iterator.next()
  87. expect(first).toMatchObject({ done: false, value: { generation: 1, value: 'first' } })
  88. if (first.done) throw new Error('fixture generation ended early')
  89. first.value.accept()
  90. const second = await iterator.next()
  91. expect(second).toMatchObject({ done: false, value: { generation: 2, value: 'second' } })
  92. if (second.done) throw new Error('fixture replacement ended early')
  93. second.value.accept()
  94. await stream.dispose()
  95. })
  96. it('permits one isolated retry while the Host remains available', async () => {
  97. const source = hostSource(true)
  98. const first = new RemoteStreamCarrierError('first carrier failure')
  99. const repeated = new RemoteStreamCarrierError('isolated retry failed')
  100. const carrierFailed = vi.fn<(error: RemoteStreamCarrierError) => void>()
  101. const stream = supervisor(source.connection, [
  102. { terminal: first },
  103. { terminal: repeated },
  104. ], carrierFailed)
  105. await expect(stream[Symbol.asyncIterator]().next()).rejects.toBe(repeated)
  106. expect(carrierFailed).toHaveBeenNthCalledWith(1, first)
  107. expect(carrierFailed).toHaveBeenNthCalledWith(2, repeated)
  108. })
  109. it('waits for a replacement Host generation after observing unavailability', async () => {
  110. let available = false
  111. let listener: (() => void) | undefined
  112. const subscribed = Promise.withResolvers<undefined>()
  113. const connection = {
  114. hostDescription: {
  115. getSnapshot: () => available ? DESCRIPTION : undefined,
  116. subscribe: (value: () => void) => {
  117. listener = value
  118. subscribed.resolve(undefined)
  119. return () => { listener = undefined }
  120. },
  121. },
  122. }
  123. let opened = 0
  124. const stream = new RemoteStream(connection, {
  125. name: 'fixture stream',
  126. open: scripted([
  127. { terminal: new RemoteStreamCarrierError('offline') },
  128. { values: ['ready'], hold: true },
  129. ], () => { opened++ }),
  130. ended: () => new Error('ended'),
  131. })
  132. const pending = stream[Symbol.asyncIterator]().next()
  133. await vi.waitFor(() => { expect(opened).toBe(1) })
  134. await subscribed.promise
  135. listener?.()
  136. expect(opened).toBe(1)
  137. available = true
  138. listener?.()
  139. await expect(pending).resolves.toMatchObject({
  140. done: false,
  141. value: { generation: 2, value: 'ready' },
  142. })
  143. await stream.dispose()
  144. })
  145. it('stops a pending retry when the logical stream is disposed', async () => {
  146. const source = hostSource(false)
  147. let opened = 0
  148. const stream = new RemoteStream(source.connection, {
  149. name: 'fixture stream',
  150. open: scripted([
  151. { terminal: new RemoteStreamCarrierError('offline') },
  152. ], () => { opened++ }),
  153. ended: () => new Error('ended'),
  154. })
  155. const pending = stream[Symbol.asyncIterator]().next()
  156. await vi.waitFor(() => { expect(opened).toBe(1) })
  157. source.publish(false)
  158. await stream.dispose()
  159. await expect(pending).resolves.toEqual({ done: true, value: undefined })
  160. })
  161. it('contains a Host publication during subscription setup', async () => {
  162. let reads = 0
  163. let disposed = 0
  164. const connection = {
  165. hostDescription: {
  166. getSnapshot: () => reads++ === 0 ? undefined : DESCRIPTION,
  167. subscribe: (listener: () => void) => {
  168. listener()
  169. return () => { disposed++ }
  170. },
  171. },
  172. }
  173. const stream = supervisor(connection, [
  174. { terminal: new RemoteStreamCarrierError('offline') },
  175. { values: ['ready'], hold: true },
  176. ])
  177. await expect(stream[Symbol.asyncIterator]().next()).resolves.toMatchObject({
  178. value: { generation: 2, value: 'ready' },
  179. })
  180. expect(disposed).toBe(1)
  181. await stream.dispose()
  182. })
  183. it('restarts with a fresh physical generation', async () => {
  184. const source = hostSource(true)
  185. const stream = supervisor(source.connection, [
  186. { values: ['first'], hold: true },
  187. { values: ['second'], hold: true },
  188. ])
  189. const iterator = stream[Symbol.asyncIterator]()
  190. await expect(iterator.next()).resolves.toMatchObject({ value: { generation: 1, value: 'first' } })
  191. stream.restart()
  192. await expect(iterator.next()).resolves.toMatchObject({ value: { generation: 2, value: 'second' } })
  193. await stream.dispose()
  194. })
  195. it('drops values and cancellation failures from a replaced generation', async () => {
  196. const source = hostSource(true)
  197. const stream = supervisor(source.connection, [
  198. { values: ['first', 'stale'] },
  199. {
  200. values: ['second'],
  201. hold: true,
  202. afterAbortError: new Error('replaced generation cancelled'),
  203. },
  204. { values: ['third'], hold: true },
  205. ])
  206. const iterator = stream[Symbol.asyncIterator]()
  207. const first = await iterator.next()
  208. if (first.done) throw new Error('fixture generation ended early')
  209. stream.restart()
  210. first.value.accept()
  211. await expect(iterator.next()).resolves.toMatchObject({
  212. value: { generation: 2, value: 'second' },
  213. })
  214. stream.restart()
  215. await expect(iterator.next()).resolves.toMatchObject({
  216. value: { generation: 3, value: 'third' },
  217. })
  218. await stream.dispose()
  219. })
  220. it('honors replacement requested by carrier diagnostics', async () => {
  221. const source = hostSource(true)
  222. const holder: { stream?: RemoteStream<string> } = {}
  223. const carrierFailed = vi.fn(() => { holder.stream?.restart() })
  224. const stream = supervisor(source.connection, [
  225. { terminal: new RemoteStreamCarrierError('replace this generation') },
  226. { values: ['ready'], hold: true },
  227. ], carrierFailed)
  228. holder.stream = stream
  229. await expect(stream[Symbol.asyncIterator]().next()).resolves.toMatchObject({
  230. value: { generation: 2, value: 'ready' },
  231. })
  232. expect(carrierFailed).toHaveBeenCalledOnce()
  233. await stream.dispose()
  234. })
  235. it('contains replacement during Host-readiness subscription setup', async () => {
  236. const holder: { stream?: RemoteStream<string> } = {}
  237. let subscriptions = 0
  238. const connection = {
  239. hostDescription: {
  240. getSnapshot: () => undefined,
  241. subscribe: () => {
  242. subscriptions++
  243. holder.stream?.restart()
  244. return () => {}
  245. },
  246. },
  247. }
  248. const stream = supervisor(connection, [
  249. { terminal: new RemoteStreamCarrierError('offline') },
  250. { values: ['ready'], hold: true },
  251. ])
  252. holder.stream = stream
  253. await expect(stream[Symbol.asyncIterator]().next()).resolves.toMatchObject({
  254. value: { generation: 2, value: 'ready' },
  255. })
  256. expect(subscriptions).toBe(1)
  257. await stream.dispose()
  258. })
  259. it('waits for generation cleanup during disposal', async () => {
  260. const source = hostSource(true)
  261. const release = Promise.withResolvers<undefined>()
  262. let closed = false
  263. const stream = supervisor(source.connection, [{
  264. values: ['ready'],
  265. hold: true,
  266. close: async () => {
  267. await release.promise
  268. closed = true
  269. },
  270. }])
  271. const iterator = stream[Symbol.asyncIterator]()
  272. await iterator.next()
  273. const pending = iterator.next()
  274. const disposing = stream.dispose()
  275. expect(stream.dispose()).toBe(disposing)
  276. await Promise.resolve()
  277. expect(closed).toBe(false)
  278. release.resolve(undefined)
  279. await expect(disposing).resolves.toBeUndefined()
  280. await expect(pending).resolves.toEqual({ done: true, value: undefined })
  281. expect(closed).toBe(true)
  282. })
  283. it('uses the domain normal-end classification and permits one consumer', async () => {
  284. const source = hostSource(true)
  285. const stream = supervisor<string>(source.connection, [{}])
  286. const iterator = stream[Symbol.asyncIterator]()
  287. expect(() => stream[Symbol.asyncIterator]()).toThrow('already has a consumer')
  288. await expect(iterator.next()).rejects.toThrow('generation ended before acceptance')
  289. await stream.dispose()
  290. })
  291. it('can be disposed before consumption and ignores later restart', async () => {
  292. const source = hostSource(true)
  293. const stream = supervisor<string>(source.connection, [])
  294. await stream.dispose()
  295. expect(stream.signal.aborted).toBe(true)
  296. stream.restart()
  297. await expect(stream[Symbol.asyncIterator]().next()).resolves.toEqual({
  298. done: true,
  299. value: undefined,
  300. })
  301. })
  302. it('drops a value that arrives after disposal begins', async () => {
  303. const source = hostSource(true)
  304. const late = Promise.withResolvers<string>()
  305. const stream = supervisor(source.connection, [{ values: [late.promise] }])
  306. const pending = stream[Symbol.asyncIterator]().next()
  307. const disposing = stream.dispose()
  308. late.resolve('late')
  309. await expect(pending).resolves.toEqual({ done: true, value: undefined })
  310. await disposing
  311. })
  312. })