1
0

reference-ownership.client.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324
  1. /** Source-labelled Client references over real history transport and scoped Contexts. */
  2. import { Context } from '@deepseek-ai/cordis'
  3. import { describe, expect, onTestFinished, vi } from 'vitest'
  4. import type {
  5. SessionReference, SessionReferenceSource, SessionRetainInfo,
  6. } from '@deepseek-ai/dsh-api-session-controller/client'
  7. import { SessionId } from '@deepseek-ai/dsh-session/types'
  8. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  9. import { ok, type RemoteMock } from '@deepseek-ai/dsh-remote-mock'
  10. import { createClientTest, webApp, type TestClient } from '@deepseek-ai/dsh-client-test-runtime/src/assembly/index.ts'
  11. import { ClientSessions } from '../src/client/sessions/service.ts'
  12. import { FOLLOW, followScript, type HistoryAnswer } from './remote/session.client.ts'
  13. declare module '@deepseek-ai/dsh-api-session-controller/client' {
  14. interface SessionReferenceSourceMap {
  15. referenceTestView: unknown
  16. referenceTestWork: unknown
  17. }
  18. }
  19. const viewSource: SessionReferenceSource = 'referenceTestView'
  20. const workSource: SessionReferenceSource = 'referenceTestWork'
  21. const ID = SessionId('reference-session')
  22. const EMPTY_HISTORY = ok({ records: [], hasMore: false })
  23. const it = createClientTest({ roster: webApp.closure(['@deepseek-ai/dsh-api-gateway']) })
  24. async function bench(mock: RemoteMock, start: () => Promise<TestClient>, listed = true) {
  25. const client = await start()
  26. const ctx = new Context()
  27. const svc = new ClientSessions(ctx, client.ctx.remote)
  28. const unblock: Array<() => void> = []
  29. onTestFinished(async () => {
  30. for (const finish of unblock) finish()
  31. await ctx.fiber.dispose()
  32. })
  33. mock.stream(FOLLOW, followScript(EMPTY_HISTORY))
  34. const feed = async (include: boolean): Promise<void> => {
  35. mock.remote.session.list.mockResolvedValue(ok({ items: include
  36. ? [{ sessionId: ID, updatedAt: 1, running: false, blank: true }]
  37. : [] }))
  38. await svc.refresh()
  39. }
  40. if (listed) await feed(true)
  41. return { svc, ctx, mock, feed, unblock }
  42. }
  43. describe('Client reference sources', () => {
  44. it('observes unknown identities without creating a scope, reference, or history request', async ({ mock, start }) => {
  45. const b = await bench(mock, start, false)
  46. const source = b.svc.retainInfo(ID)
  47. const before = source.getSnapshot()
  48. const stop = source.subscribe(vi.fn())
  49. expect(b.svc.retainInfo(ID)).toBe(source)
  50. expect(source.getSnapshot()).toBe(before)
  51. expect(before).toEqual({ referenceCount: 0, retainedBy: {} })
  52. expect('set' in source).toBe(false)
  53. expect(b.svc.scope(ID)).toBeUndefined()
  54. expect(b.svc.binding(ID)).toBeUndefined()
  55. expect(mock.log.requests(FOLLOW)).toHaveLength(0)
  56. expect(() => b.svc.retain(ID, { source: viewSource })).toThrow('unknown session')
  57. expect(source.getSnapshot()).toBe(before)
  58. stop()
  59. })
  60. it('shares initial opening while exposing independent source contributions before the await', async ({ mock, start }) => {
  61. const b = await bench(mock, start)
  62. const opening = Promise.withResolvers<Awaited<HistoryAnswer>>()
  63. const entered = Promise.withResolvers<undefined>()
  64. b.unblock.push(() => { opening.resolve(EMPTY_HISTORY) })
  65. mock.stream(FOLLOW, followScript(() => { entered.resolve(undefined); return opening.promise }))
  66. const source = b.svc.retainInfo(ID)
  67. const first = b.svc.retain(ID, { source: viewSource })
  68. const binding = b.svc.binding(ID)
  69. const second = b.svc.retain(ID, { source: workSource })
  70. const acquisitions = Promise.allSettled([first.ready, second.ready])
  71. expect(binding).toBeDefined()
  72. expect(source.getSnapshot()).toEqual({ referenceCount: 2, retainedBy: { referenceTestView: 1, referenceTestWork: 1 } })
  73. expect(b.svc.list.getSnapshot().byId[ID]?.retainedBy).toBe(source.getSnapshot().retainedBy)
  74. await entered.promise
  75. expect(mock.log.requests(FOLLOW)).toHaveLength(1)
  76. opening.resolve(EMPTY_HISTORY)
  77. await acquisitions
  78. using a = first
  79. using c = second
  80. expect(a.binding).toBe(binding)
  81. expect(c.binding).toBe(binding)
  82. expect(a.binding.session.getSnapshot().openState).toBe('open')
  83. a.release()
  84. a[Symbol.dispose]()
  85. expect(source.getSnapshot()).toEqual({ referenceCount: 1, retainedBy: { referenceTestWork: 1 } })
  86. expect(() => a.binding).toThrow('is released')
  87. c.release()
  88. expect(source.getSnapshot()).toEqual({ referenceCount: 0, retainedBy: {} })
  89. expect(b.svc.binding(ID)).toBeUndefined()
  90. })
  91. it('counts repeated uses of one source and omits it after the final release', async ({ mock, start }) => {
  92. const b = await bench(mock, start)
  93. using first = b.svc.retain(ID, { source: workSource, signal: undefined })
  94. using second = b.svc.retain(ID, { source: workSource })
  95. await Promise.all([first.ready, second.ready])
  96. const source = b.svc.retainInfo(ID)
  97. expect(source.getSnapshot()).toEqual({ referenceCount: 2, retainedBy: { referenceTestWork: 2 } })
  98. first.release()
  99. expect(source.getSnapshot().retainedBy.referenceTestWork).toBe(1)
  100. second.release()
  101. expect(source.getSnapshot().retainedBy.referenceTestWork).toBeUndefined()
  102. expect(Object.keys(source.getSnapshot().retainedBy)).toEqual([])
  103. })
  104. it('keeps counts through catalog refresh/removal and projects them when the row returns', async ({ mock, start }) => {
  105. const b = await bench(mock, start)
  106. using reference = b.svc.retain(ID, { source: viewSource })
  107. await reference.ready
  108. const binding = reference.binding
  109. const source = b.svc.retainInfo(ID)
  110. const snapshot = source.getSnapshot()
  111. await b.feed(true)
  112. expect(b.svc.list.getSnapshot().byId[ID]?.retainedBy).toBe(snapshot.retainedBy)
  113. b.svc.handleSessionRemoved(ID)
  114. await vi.waitFor(() => {
  115. expect(b.svc.list.getSnapshot().byId[ID]?.retainedBy).toBe(snapshot.retainedBy)
  116. })
  117. expect(b.svc.list.getSnapshot().ids).not.toContain(ID)
  118. expect(source.getSnapshot()).toBe(snapshot)
  119. expect(b.svc.binding(ID)).toBe(binding)
  120. await b.feed(true)
  121. expect(b.svc.list.getSnapshot().byId[ID]?.retainedBy).toBe(snapshot.retainedBy)
  122. reference.release()
  123. expect(b.svc.list.getSnapshot().byId[ID]?.retainedBy).toEqual({})
  124. })
  125. it('retains a synchronous Gateway Context before catalog discovery without history I/O', async ({ mock, start }) => {
  126. const b = await bench(mock, start, false)
  127. const source = b.svc.retainInfo(ID)
  128. using reference = b.svc.retainAgentScope(ID)
  129. expect(reference.binding.ctx).toBe(b.svc.scope(ID))
  130. expect(reference.binding.session.getSnapshot().openState).toBe('cold')
  131. expect(source.getSnapshot()).toEqual({ referenceCount: 1, retainedBy: { gateway: 1 } })
  132. expect(b.svc.list.getSnapshot().ids).not.toContain(ID)
  133. expect(b.svc.list.getSnapshot().byId[ID]?.retainedBy).toEqual({ gateway: 1 })
  134. expect(mock.log.requests(FOLLOW)).toHaveLength(0)
  135. expect(mock.remote.subagents.list).not.toHaveBeenCalled()
  136. await b.feed(true)
  137. expect(b.svc.list.getSnapshot().byId[ID]?.retainedBy).toEqual({ gateway: 1 })
  138. })
  139. for (const kind of ['remote', 'unexpected'] as const) it(`settles ${kind} initial opening failure through Session state and allows another acquisition after release`, async ({ mock, start }) => {
  140. const b = await bench(mock, start)
  141. using local = b.svc.retainAgentScope(ID)
  142. const binding = local.binding
  143. const failure = kind === 'remote'
  144. ? new RemoteError('session/not-found', 'opening failed', { sessionId: ID })
  145. : new Error('opening failed')
  146. mock.stream(FOLLOW, followScript(() => Promise.reject(failure)))
  147. const failed = b.svc.retain(ID, { source: viewSource })
  148. await expect(failed.ready).resolves.toBe(binding)
  149. expect(b.svc.retainInfo(ID).getSnapshot()).toEqual({ referenceCount: 2, retainedBy: { gateway: 1, referenceTestView: 1 } })
  150. expect(binding.session.getSnapshot().openState).toBe('error')
  151. failed.release()
  152. mock.stream(FOLLOW, followScript(EMPTY_HISTORY))
  153. using retried = b.svc.retain(ID, { source: workSource })
  154. await retried.ready
  155. expect(retried.binding).toBe(binding)
  156. expect(retried.binding.session.getSnapshot().openState).toBe('open')
  157. })
  158. it('cancels only one waiter while another owns the shared initial opening', async ({ mock, start }) => {
  159. const b = await bench(mock, start)
  160. const opening = Promise.withResolvers<Awaited<HistoryAnswer>>()
  161. b.unblock.push(() => { opening.resolve(EMPTY_HISTORY) })
  162. mock.stream(FOLLOW, followScript(opening.promise))
  163. const controller = new AbortController()
  164. const cancelled = b.svc.retain(ID, { source: viewSource, signal: controller.signal })
  165. const survivor = b.svc.retain(ID, { source: workSource })
  166. const reason = new Error('waiter cancelled')
  167. const rejected = expect(cancelled.ready).rejects.toBe(reason)
  168. controller.abort(reason)
  169. await rejected
  170. expect(b.svc.retainInfo(ID).getSnapshot()).toEqual({ referenceCount: 2, retainedBy: { referenceTestView: 1, referenceTestWork: 1 } })
  171. cancelled.release()
  172. expect(b.svc.retainInfo(ID).getSnapshot()).toEqual({ referenceCount: 1, retainedBy: { referenceTestWork: 1 } })
  173. opening.resolve(EMPTY_HISTORY)
  174. using reference = survivor
  175. await reference.ready
  176. expect(reference.binding.session.getSnapshot().openState).toBe('open')
  177. expect(mock.log.requests(FOLLOW)).toHaveLength(1)
  178. })
  179. it('rejects a previously cancelled acquisition without creating a generation or publishing counts', async ({ mock, start }) => {
  180. const b = await bench(mock, start)
  181. const controller = new AbortController()
  182. const reason = new Error('acquisition cancelled')
  183. controller.abort(reason)
  184. const changed = vi.fn()
  185. const source = b.svc.retainInfo(ID)
  186. b.ctx.effect(() => source.subscribe(changed), 'test: reference observation')
  187. const before = source.getSnapshot()
  188. expect(() => b.svc.retain(ID, { source: viewSource, signal: controller.signal })).toThrow(reason)
  189. expect(source.getSnapshot()).toBe(before)
  190. expect(changed).not.toHaveBeenCalled()
  191. expect(b.svc.binding(ID)).toBeUndefined()
  192. expect(mock.log.requests(FOLLOW)).toHaveLength(0)
  193. })
  194. it('releases a synchronously cancelled readiness wait through using()', async ({ mock, start }) => {
  195. const b = await bench(mock, start)
  196. const controller = new AbortController()
  197. const failure = new Error('owner ended during acquisition')
  198. const source = b.svc.retainInfo(ID)
  199. b.ctx.effect(() => source.subscribe(() => {
  200. if (source.getSnapshot().referenceCount > 0) controller.abort(failure)
  201. }), 'test: acquisition cancellation')
  202. await expect(b.svc.using(ID, { source: viewSource, signal: controller.signal }, () => undefined)).rejects.toBe(failure)
  203. expect(source.getSnapshot()).toEqual({ referenceCount: 0, retainedBy: {} })
  204. expect(b.svc.binding(ID)).toBeUndefined()
  205. })
  206. it('withdraws the old generation before an observer retains a same-id replacement', async ({ mock, start }) => {
  207. const b = await bench(mock, start)
  208. using old = b.svc.retain(ID, { source: viewSource })
  209. await old.ready
  210. const binding = old.binding
  211. const source = b.svc.retainInfo(ID)
  212. const replacement = Promise.withResolvers<SessionReference>()
  213. const stop = source.subscribe(() => {
  214. if (source.getSnapshot().referenceCount !== 0) return
  215. stop()
  216. replacement.resolve(b.svc.retainAgentScope(ID))
  217. })
  218. b.ctx.effect(() => stop, 'test: generation replacement')
  219. old.release()
  220. using reference = await replacement.promise
  221. await binding.ctx.fiber.dispose()
  222. old.release()
  223. expect(reference.binding).not.toBe(binding)
  224. expect(reference.binding.session.getSnapshot().openState).toBe('cold')
  225. expect(b.svc.sessionOf(binding.ctx)).toBeUndefined()
  226. expect(source.getSnapshot()).toEqual({ referenceCount: 1, retainedBy: { gateway: 1 } })
  227. expect(mock.log.requests(FOLLOW)).toHaveLength(1)
  228. })
  229. it('keeps one observable across replacement while late old cleanup cannot change its counts', async ({ mock, start }) => {
  230. const b = await bench(mock, start)
  231. const source = b.svc.retainInfo(ID)
  232. const old = b.svc.retain(ID, { source: viewSource })
  233. await old.ready
  234. const binding = old.binding
  235. const entered = Promise.withResolvers<undefined>()
  236. const resume = Promise.withResolvers<undefined>()
  237. b.unblock.push(() => { resume.resolve(undefined) })
  238. binding.ctx.effect(() => async () => { entered.resolve(undefined); await resume.promise }, 'test: delayed generation cleanup')
  239. old.release()
  240. expect(b.svc.binding(ID)).toBeUndefined()
  241. expect(source.getSnapshot().referenceCount).toBe(0)
  242. await entered.promise
  243. using replacement = b.svc.retain(ID, { source: workSource })
  244. await replacement.ready
  245. expect(replacement.binding).not.toBe(binding)
  246. const snapshot = source.getSnapshot()
  247. resume.resolve(undefined)
  248. await binding.ctx.fiber.dispose()
  249. old.release()
  250. expect(b.svc.retainInfo(ID)).toBe(source)
  251. expect(source.getSnapshot()).toBe(snapshot)
  252. expect(b.svc.sessionOf(binding.ctx)).toBeUndefined()
  253. expect(b.svc.binding(ID)).toBe(replacement.binding)
  254. })
  255. it('invalidates references and zeroes their observable on root disposal', async ({ mock, start }) => {
  256. const b = await bench(mock, start)
  257. const reference = b.svc.retain(ID, { source: viewSource })
  258. await reference.ready
  259. const source = b.svc.retainInfo(ID)
  260. await b.ctx.fiber.dispose()
  261. expect(source.getSnapshot()).toEqual({ referenceCount: 0, retainedBy: {} })
  262. expect(() => reference.binding).toThrow('is released')
  263. reference.release()
  264. expect(() => b.svc.retain(ID, { source: workSource })).toThrow('Controller is disposed')
  265. expect(() => b.svc.retainAgentScope(ID)).toThrow('Controller is disposed')
  266. })
  267. })
  268. describe('ClientSessions.using', () => {
  269. it('awaits callback settlement before releasing and returns its result', async ({ mock, start }) => {
  270. const b = await bench(mock, start)
  271. const entered = Promise.withResolvers<SessionReference>()
  272. const response = Promise.withResolvers<number>()
  273. b.unblock.push(() => { response.resolve(7) })
  274. const using = b.svc.using(ID, { source: workSource, signal: undefined }, (reference) => {
  275. entered.resolve(reference)
  276. return response.promise
  277. })
  278. const reference = await entered.promise
  279. expect(b.svc.retainInfo(ID).getSnapshot().referenceCount).toBe(1)
  280. response.resolve(7)
  281. await expect(using).resolves.toBe(7)
  282. expect(b.svc.retainInfo(ID).getSnapshot().referenceCount).toBe(0)
  283. expect(() => reference.binding).toThrow('is released')
  284. })
  285. for (const kind of ['sync', 'async'] as const) it(`releases and propagates a ${kind} callback failure`, async ({ mock, start }) => {
  286. const b = await bench(mock, start)
  287. const failure = new Error('callback failed')
  288. await expect(b.svc.using(ID, { source: workSource }, () => {
  289. if (kind === 'sync') throw failure
  290. return Promise.reject(failure)
  291. })).rejects.toBe(failure)
  292. expect(b.svc.retainInfo(ID).getSnapshot()).toEqual({ referenceCount: 0, retainedBy: {} })
  293. })
  294. it('calls the operation after a stateful opening failure', async ({ mock, start }) => {
  295. const b = await bench(mock, start)
  296. const operation = vi.fn(() => 7)
  297. mock.stream(FOLLOW, followScript(() => Promise.reject(new Error('cannot open'))))
  298. await expect(b.svc.using(ID, { source: workSource }, operation)).resolves.toBe(7)
  299. expect(operation).toHaveBeenCalledOnce()
  300. const info: SessionRetainInfo = b.svc.retainInfo(ID).getSnapshot()
  301. expect(info).toEqual({ referenceCount: 0, retainedBy: {} })
  302. expect(b.svc.binding(ID)).toBeUndefined()
  303. })
  304. })