change-feed.client.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277
  1. /**
  2. * The change feed's promises: one Host stream per session, delivery by absolute
  3. * path, and a follower's life bounded by its signal or by
  4. * the stream's end.
  5. */
  6. import type { SessionId } from '@deepseek-ai/dsh-session/types'
  7. import { describe, expect, it } from 'vitest'
  8. import { ChangeFeed } from '../src/client/change-feed.ts'
  9. import type { WorkspaceFileWatchFrame } from '../src/types.ts'
  10. import { FakeRemote, peek, settle } from './fake-remote.client.ts'
  11. const S1 = 's1' as SessionId
  12. const S2 = 's2' as SessionId
  13. function harness() {
  14. const remote = new FakeRemote()
  15. const feed = new ChangeFeed(remote)
  16. const follow = (sessionId: SessionId, path: string, controller = new AbortController()) => {
  17. const follower = feed.follow(sessionId, controller.signal)
  18. follower.bind(path)
  19. return { it: follower[Symbol.asyncIterator](), controller }
  20. }
  21. return { remote, feed, follow }
  22. }
  23. describe('ChangeFeed — one Host stream per session', () => {
  24. it('starts a later follower from the existing session acknowledgement without opening another stream', async () => {
  25. const { remote, feed } = harness()
  26. const controller = new AbortController()
  27. const first = feed.follow(S1, controller.signal)
  28. try {
  29. await expect(first.ready).resolves.toBe(true)
  30. const second = feed.follow(S1, controller.signal)
  31. await expect(second.ready).resolves.toBe(true)
  32. expect(remote.calls).toEqual(['changes', 'accept'])
  33. expect(remote.opened).toHaveLength(1)
  34. second.bind('/w/second.txt')
  35. const iterator = second[Symbol.asyncIterator]()
  36. await remote.opened[0]!.source.deliver({
  37. kind: 'change', change: { absolutePath: '/w/second.txt', version: 'v1' },
  38. })
  39. await expect(iterator.next()).resolves.toEqual({ done: false, value: { kind: 'changed', version: 'v1' } })
  40. controller.abort()
  41. await expect(iterator.next()).resolves.toEqual({ done: true, value: undefined })
  42. await feed.settle()
  43. expect(remote.disposed).toEqual(['workspace file changes of s1'])
  44. } finally {
  45. controller.abort()
  46. await feed.settle()
  47. }
  48. })
  49. it('shares one stream among the followers of a session and opens another per session', async () => {
  50. const { remote, follow } = harness()
  51. follow(S1, '/w/a.txt')
  52. follow(S1, '/w/b.txt')
  53. await settle()
  54. expect(remote.opened.map(o => o.sessionId)).toEqual([S1])
  55. follow(S2, '/w/a.txt')
  56. await settle()
  57. expect(remote.opened.map(o => o.sessionId)).toEqual([S1, S2])
  58. })
  59. it('disposes the session stream when its last follower leaves and reopens for the next', async () => {
  60. const { remote, follow } = harness()
  61. const a = follow(S1, '/w/a.txt')
  62. const b = follow(S1, '/w/b.txt')
  63. await settle()
  64. a.controller.abort()
  65. await settle()
  66. expect(remote.disposed).toEqual([])
  67. b.controller.abort()
  68. await settle()
  69. expect(remote.disposed).toEqual(['workspace file changes of s1'])
  70. expect(remote.opened[0]!.source.aborted).toBe(true)
  71. follow(S1, '/w/c.txt')
  72. await settle()
  73. expect(remote.opened).toHaveLength(2)
  74. })
  75. it('opens the next stream of a session only after the previous dispose settled, and settle() waits for it', async () => {
  76. const { remote, feed, follow } = harness()
  77. let release!: () => void
  78. remote.disposeGate = new Promise<void>((resolve) => { release = resolve })
  79. const a = follow(S1, '/w/a.txt')
  80. await settle()
  81. a.controller.abort()
  82. await settle()
  83. expect(remote.disposed).toEqual(['workspace file changes of s1'])
  84. // The next follower registers at once, but its Host stream waits for the close.
  85. follow(S1, '/w/b.txt')
  86. await settle()
  87. expect(remote.opened).toHaveLength(1)
  88. let settled = false
  89. void feed.settle().then(() => { settled = true })
  90. await settle()
  91. expect(settled).toBe(false)
  92. release()
  93. await settle()
  94. expect(remote.opened).toHaveLength(2)
  95. expect(settled).toBe(true)
  96. // Nothing left closing: settle() resolves at once.
  97. await feed.settle()
  98. })
  99. it('keeps waiting for the newest close when two closes of one session overlap', async () => {
  100. const { remote, feed, follow } = harness()
  101. let release!: () => void
  102. remote.disposeGate = new Promise<void>((resolve) => { release = resolve })
  103. const a = follow(S1, '/w/a.txt')
  104. await settle()
  105. a.controller.abort()
  106. await settle()
  107. // The second feed waits for the first close, then its only follower leaves too.
  108. const b = follow(S1, '/w/b.txt')
  109. await settle()
  110. b.controller.abort()
  111. await settle()
  112. let settled = false
  113. void feed.settle().then(() => { settled = true })
  114. release()
  115. await settle()
  116. await settle()
  117. expect(settled).toBe(true)
  118. expect(remote.disposed).toHaveLength(2)
  119. follow(S1, '/w/c.txt')
  120. await settle()
  121. expect(remote.opened).toHaveLength(3)
  122. })
  123. it('treats a dispose that rejects as settled, so the next stream still opens', async () => {
  124. const { remote, feed, follow } = harness()
  125. // Rejected only once dispose() has taken the gate, so the rejection always has a handler.
  126. let fail!: (error: Error) => void
  127. remote.disposeGate = new Promise<void>((_resolve, reject) => { fail = reject })
  128. const a = follow(S1, '/w/a.txt')
  129. await settle()
  130. a.controller.abort()
  131. await settle()
  132. follow(S1, '/w/b.txt')
  133. await settle()
  134. expect(remote.opened).toHaveLength(1)
  135. fail(new Error('carrier gone'))
  136. await settle()
  137. expect(remote.opened).toHaveLength(2)
  138. await feed.settle()
  139. })
  140. })
  141. describe('ChangeFeed — delivery', () => {
  142. it('routes a frame to the followers of its path, whichever separator the Host spells', async () => {
  143. const { remote, follow } = harness()
  144. // A stat and a change frame may spell the same Host path with different separators.
  145. const mine = follow(S1, 'C:/w/a b.txt')
  146. const twin = follow(S1, 'C:/w/a b.txt')
  147. const other = follow(S1, 'C:/w/other.txt')
  148. await settle()
  149. const source = remote.opened[0]!.source
  150. source.push({ kind: 'change', change: { absolutePath: 'C:/w/a b.txt', version: 'v1' } })
  151. source.push({ kind: 'change', change: { absolutePath: 'C:\\w\\a b.txt', absent: true } })
  152. await expect(mine.it.next()).resolves.toEqual({ done: false, value: { kind: 'changed', version: 'v1' } })
  153. await expect(mine.it.next()).resolves.toEqual({ done: false, value: { kind: 'absent' } })
  154. await expect(twin.it.next()).resolves.toEqual({ done: false, value: { kind: 'changed', version: 'v1' } })
  155. await expect(peek(other.it)).resolves.toBe('silent')
  156. // One follower of a path leaving does not silence the other.
  157. mine.controller.abort()
  158. await settle()
  159. source.push({ kind: 'change', change: { absolutePath: 'C:/w/a b.txt', version: 'v2' } })
  160. await expect(twin.it.next()).resolves.toEqual({ done: false, value: { kind: 'absent' } })
  161. await expect(twin.it.next()).resolves.toEqual({ done: false, value: { kind: 'changed', version: 'v2' } })
  162. })
  163. it('queues frames reported before the consumer starts pulling', async () => {
  164. const { remote, follow } = harness()
  165. const mine = follow(S1, '/w/a.txt')
  166. await settle()
  167. remote.opened[0]!.source.push({ kind: 'change', change: { absolutePath: '/w/a.txt', version: 'v1' } })
  168. await settle()
  169. await expect(mine.it.next()).resolves.toEqual({ done: false, value: { kind: 'changed', version: 'v1' } })
  170. })
  171. it('keeps sessions apart', async () => {
  172. const { remote, follow } = harness()
  173. const one = follow(S1, '/w/a.txt')
  174. const two = follow(S2, '/w/a.txt')
  175. await settle()
  176. remote.opened[1]!.source.push({ kind: 'change', change: { absolutePath: '/w/a.txt', version: 'v2' } })
  177. await expect(two.it.next()).resolves.toEqual({ done: false, value: { kind: 'changed', version: 'v2' } })
  178. await expect(peek(one.it)).resolves.toBe('silent')
  179. })
  180. })
  181. describe('ChangeFeed — a follower ends', () => {
  182. it('ends every follower and disposes the session stream for an unknown wire frame kind', async () => {
  183. const { remote, feed } = harness()
  184. const controller = new AbortController()
  185. const first = feed.follow(S1, controller.signal)
  186. const second = feed.follow(S1, controller.signal)
  187. const firstIterator = first[Symbol.asyncIterator]()
  188. const secondIterator = second[Symbol.asyncIterator]()
  189. try {
  190. await expect(Promise.all([first.ready, second.ready])).resolves.toEqual([true, true])
  191. const endings = Promise.all([firstIterator.next(), secondIterator.next()])
  192. const source = remote.opened[0]!.source
  193. // The Remote double supplies decoded wire data, including an unknown protocol tag.
  194. const wireFrame: unknown = JSON.parse('{"kind":"future-frame"}')
  195. source.push(wireFrame as WorkspaceFileWatchFrame)
  196. await expect(endings).resolves.toEqual([
  197. { done: true, value: undefined },
  198. { done: true, value: undefined },
  199. ])
  200. await feed.settle()
  201. expect(source.aborted).toBe(true)
  202. expect(remote.disposed).toEqual(['workspace file changes of s1'])
  203. expect(remote.opened).toHaveLength(1)
  204. } finally {
  205. controller.abort()
  206. await Promise.all([firstIterator.return?.(), secondIterator.return?.()])
  207. await feed.settle()
  208. }
  209. })
  210. it('ends on its signal and drops nothing queued before it', async () => {
  211. const { remote, follow } = harness()
  212. const mine = follow(S1, '/w/a.txt')
  213. await settle()
  214. remote.opened[0]!.source.push({ kind: 'change', change: { absolutePath: '/w/a.txt', version: 'v1' } })
  215. await settle()
  216. mine.controller.abort()
  217. await expect(mine.it.next()).resolves.toEqual({ done: false, value: { kind: 'changed', version: 'v1' } })
  218. await expect(mine.it.next()).resolves.toEqual({ done: true, value: undefined })
  219. })
  220. it('is empty when the signal is already aborted, without opening a stream', async () => {
  221. const { remote, feed } = harness()
  222. const controller = new AbortController()
  223. controller.abort()
  224. const it = feed.follow(S1, controller.signal)[Symbol.asyncIterator]()
  225. await expect(it.next()).resolves.toEqual({ done: true, value: undefined })
  226. await settle()
  227. expect(remote.opened).toEqual([])
  228. })
  229. it('unregisters when the consumer breaks out after a notice', async () => {
  230. const { remote, follow } = harness()
  231. const mine = follow(S1, '/w/a.txt')
  232. await settle()
  233. remote.opened[0]!.source.push({ kind: 'change', change: { absolutePath: '/w/a.txt', version: 'v1' } })
  234. await mine.it.next()
  235. await mine.it.return?.()
  236. await settle()
  237. expect(remote.disposed).toEqual(['workspace file changes of s1'])
  238. })
  239. it('ends every follower when the Host closes the session stream', async () => {
  240. const { remote, follow } = harness()
  241. const a = follow(S1, '/w/a.txt')
  242. const b = follow(S1, '/w/b.txt')
  243. await settle()
  244. remote.opened[0]!.source.end()
  245. await expect(a.it.next()).resolves.toEqual({ done: true, value: undefined })
  246. await expect(b.it.next()).resolves.toEqual({ done: true, value: undefined })
  247. // The next follower starts a fresh stream rather than joining the dead one.
  248. follow(S1, '/w/c.txt')
  249. await settle()
  250. expect(remote.opened).toHaveLength(2)
  251. })
  252. it('ends every follower when the session stream fails', async () => {
  253. const { remote, follow } = harness()
  254. const a = follow(S1, '/w/a.txt')
  255. await settle()
  256. remote.opened[0]!.source.fail(new Error('carrier gone for good'))
  257. await expect(a.it.next()).resolves.toEqual({ done: true, value: undefined })
  258. })
  259. })