control-retry.client.spec.ts 13 KB

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