resources.spec.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357
  1. import { Context } from '@deepseek-ai/cordis'
  2. import AgentRegistry from '@deepseek-ai/dsh-agent'
  3. import type { Agent } from '@deepseek-ai/dsh-agent'
  4. import { unsupportedInbox } from '@deepseek-ai/dsh-agent-loop-testkit'
  5. import { Session, SessionId } from '@deepseek-ai/dsh-session'
  6. import { afterEach, describe, expect, it, vi } from 'vitest'
  7. import { SessionResources } from '../src/index.ts'
  8. const contexts: Context[] = []
  9. async function fixture() {
  10. const ctx = new Context()
  11. contexts.push(ctx)
  12. await ctx.plugin(AgentRegistry)
  13. async function owner(id: string) {
  14. const fiber = ctx.plugin(() => {})
  15. const session = Session.create(SessionId(id))
  16. const agent: Agent = {
  17. id: session.id, session, ctx: fiber.ctx, options: {}, status: 'idle',
  18. inbox: unsupportedInbox(), send() {}, followup() {}, inject() {}, cancel() {},
  19. steer: () => ({ outcome: Promise.resolve({ status: 'rejected' as const }) }),
  20. runMaintenance: task => task(new AbortController().signal),
  21. whenIdle: () => Promise.resolve(undefined),
  22. }
  23. const unregister = ctx.agents.register(agent)
  24. await unregister
  25. return { agent, async dispose() { await fiber.dispose(); await unregister() } }
  26. }
  27. return { ctx, owner }
  28. }
  29. afterEach(async () => {
  30. await Promise.all(contexts.splice(0).map(ctx => ctx.fiber.dispose()))
  31. })
  32. describe('Session browser resource ownership', () => {
  33. it('acquires once across concurrent requests and gives another Session a different resource', async () => {
  34. const { ctx, owner } = await fixture()
  35. const a = await owner('a')
  36. const b = await owner('b')
  37. const close = vi.fn(async () => {})
  38. const open = vi.fn(async (agent: Agent) => ({ value: { id: agent.id }, close }))
  39. const resources = new SessionResources(ctx, { label: 'test', exclusive: false, open })
  40. const [first, again, other] = await Promise.all([resources.get(a.agent), resources.get(a.agent), resources.get(b.agent)])
  41. expect(first).toBe(again)
  42. expect(other).not.toBe(first)
  43. expect(open).toHaveBeenCalledTimes(2)
  44. await a.dispose()
  45. expect(close).toHaveBeenCalledTimes(1)
  46. await expect(resources.get(a.agent)).rejects.toThrow('not a live browser owner')
  47. const resumed = await owner('a')
  48. expect(await resources.get(resumed.agent)).not.toBe(first)
  49. await resources.dispose()
  50. expect(close).toHaveBeenCalledTimes(3)
  51. expect(resources.available(b.agent)).toBe(false)
  52. await expect(resources.get(b.agent)).rejects.toThrow('not a live browser owner')
  53. })
  54. it('reserves an attached browser while acquisition or cleanup is pending', async () => {
  55. const { ctx, owner } = await fixture()
  56. const a = await owner('a')
  57. const b = await owner('b')
  58. const opened = Promise.withResolvers<undefined>()
  59. const released = Promise.withResolvers<undefined>()
  60. const closing = Promise.withResolvers<undefined>()
  61. const resources = new SessionResources(ctx, {
  62. label: 'attached', exclusive: true,
  63. async open() {
  64. await opened.promise
  65. return { value: {}, async close() { closing.resolve(undefined); await released.promise } }
  66. },
  67. })
  68. expect(resources.available(a.agent)).toBe(true)
  69. expect(resources.available(b.agent)).toBe(true)
  70. const first = resources.get(a.agent)
  71. expect(resources.available(a.agent)).toBe(true)
  72. expect(resources.available(b.agent)).toBe(false)
  73. await expect(resources.get(b.agent)).rejects.toThrow('already reserved')
  74. opened.resolve(undefined)
  75. await first
  76. const disposing = a.dispose()
  77. await closing.promise
  78. expect(resources.available(a.agent)).toBe(false)
  79. expect(resources.available(b.agent)).toBe(false)
  80. await expect(resources.get(b.agent)).rejects.toThrow('already reserved')
  81. released.resolve(undefined)
  82. await disposing
  83. expect(resources.available(b.agent)).toBe(true)
  84. await resources.get(b.agent)
  85. await resources.dispose()
  86. })
  87. it('releases a failed acquisition and can acquire for another Session', async () => {
  88. const { ctx, owner } = await fixture()
  89. const a = await owner('a')
  90. const b = await owner('b')
  91. const open = vi.fn().mockRejectedValueOnce(new Error('browser unavailable')).mockResolvedValue({ value: 1, close: async () => {} })
  92. const resources = new SessionResources<number>(ctx, { label: 'test', exclusive: true, open })
  93. await expect(resources.get(a.agent, new AbortController().signal)).rejects.toThrow('browser unavailable')
  94. expect(await resources.get(b.agent)).toBe(1)
  95. await resources.dispose()
  96. })
  97. it('retries failed acquisition for the same live owner without duplicating cleanup', async () => {
  98. const { ctx, owner } = await fixture()
  99. const a = await owner('a')
  100. const close = vi.fn(async () => {})
  101. const open = vi.fn().mockRejectedValueOnce(new Error('launch failed')).mockResolvedValue({ value: 1, close })
  102. const resources = new SessionResources<number>(ctx, { label: 'test', exclusive: false, open })
  103. await expect(resources.get(a.agent)).rejects.toThrow('launch failed')
  104. expect(await resources.get(a.agent)).toBe(1)
  105. await a.dispose()
  106. await resources.dispose()
  107. expect(close).toHaveBeenCalledTimes(1)
  108. })
  109. it('quiesces disposal when an in-flight acquisition fails during rollback', async () => {
  110. const { ctx, owner } = await fixture()
  111. const a = await owner('a')
  112. const entered = Promise.withResolvers<undefined>()
  113. const failed = Promise.withResolvers<never>()
  114. const resources = new SessionResources(ctx, {
  115. label: 'test', exclusive: false,
  116. async open() { entered.resolve(undefined); return failed.promise },
  117. })
  118. const acquiring = resources.get(a.agent)
  119. const rejected = expect(acquiring).rejects.toThrow('launch failed')
  120. await entered.promise
  121. const disposing = resources.dispose()
  122. failed.reject(new Error('launch failed'))
  123. await Promise.all([rejected, disposing])
  124. })
  125. it('serializes one Session while another proceeds and skips cancelled queued work', async () => {
  126. const { ctx, owner } = await fixture()
  127. const a = await owner('a')
  128. const b = await owner('b')
  129. const resources = new SessionResources(ctx, {
  130. label: 'test', exclusive: false,
  131. open: async () => ({ value: {}, close: async () => {} }),
  132. })
  133. const entered = Promise.withResolvers<undefined>()
  134. const release = Promise.withResolvers<undefined>()
  135. const first = resources.run(a.agent, new AbortController().signal, async () => {
  136. entered.resolve(undefined)
  137. await release.promise
  138. return 1
  139. })
  140. await entered.promise
  141. const abort = new AbortController()
  142. const queued = vi.fn(async () => 2)
  143. const second = resources.run(a.agent, abort.signal, queued)
  144. const rejected = expect(second).rejects.toThrow('cancel queued')
  145. abort.abort(new Error('cancel queued'))
  146. expect(await resources.run(b.agent, new AbortController().signal, async () => 3)).toBe(3)
  147. release.resolve(undefined)
  148. expect(await first).toBe(1)
  149. await rejected
  150. expect(queued).not.toHaveBeenCalled()
  151. expect(await resources.run(a.agent, new AbortController().signal, async () => 4)).toBe(4)
  152. await resources.dispose()
  153. })
  154. it.each(['get', 'run'] as const)('cancels one %s caller while another waiter retains the same acquisition', async (kind) => {
  155. const { ctx, owner } = await fixture()
  156. const a = await owner('a')
  157. const b = await owner('b')
  158. const entered = Promise.withResolvers<undefined>()
  159. const release = Promise.withResolvers<undefined>()
  160. const close = vi.fn(async () => {})
  161. let acquisitionSignal: AbortSignal | undefined
  162. const open = vi.fn(async (_agent: Agent, signal: AbortSignal) => {
  163. acquisitionSignal = signal
  164. entered.resolve(undefined)
  165. await release.promise
  166. return { value: 1, close }
  167. })
  168. const resources = new SessionResources(ctx, { label: 'test', exclusive: true, open })
  169. const controller = new AbortController()
  170. const execute = vi.fn(async (value: number) => value)
  171. const acquiring = kind === 'get'
  172. ? resources.get(a.agent, controller.signal)
  173. : resources.run(a.agent, controller.signal, execute)
  174. const canceled = acquiring.catch((error: unknown) => error)
  175. const retained = resources.get(a.agent, new AbortController().signal).catch((error: unknown) => error)
  176. try {
  177. await entered.promise
  178. controller.abort(new Error('cancel turn'))
  179. expect(await canceled).toMatchObject({ message: 'cancel turn' })
  180. expect(acquisitionSignal?.aborted).toBe(false)
  181. expect(resources.available(b.agent)).toBe(false)
  182. expect(execute).not.toHaveBeenCalled()
  183. release.resolve(undefined)
  184. expect(await retained).toBe(1)
  185. expect(await resources.get(a.agent)).toBe(1)
  186. expect(open).toHaveBeenCalledOnce()
  187. } finally {
  188. release.resolve(undefined)
  189. await resources.dispose()
  190. }
  191. expect(close).toHaveBeenCalledOnce()
  192. })
  193. it('honors caller cancellation between shared readiness and the waiting continuation', async () => {
  194. const { ctx, owner } = await fixture()
  195. const a = await owner('a')
  196. const ready = Promise.withResolvers<undefined>()
  197. const resources = new SessionResources(ctx, {
  198. label: 'test', exclusive: false,
  199. async open() { await ready.promise; return { value: 1, close: async () => {} } },
  200. })
  201. const controller = new AbortController()
  202. const error = new Error('cancel ready waiter')
  203. const retained = resources.get(a.agent)
  204. const canceled = resources.get(a.agent, controller.signal).catch((failure: unknown) => failure)
  205. const abort = retained.then(() => { controller.abort(error) })
  206. try {
  207. ready.resolve(undefined)
  208. await abort
  209. expect(await canceled).toBe(error)
  210. expect(await retained).toBe(1)
  211. expect(await resources.get(a.agent)).toBe(1)
  212. } finally {
  213. ready.resolve(undefined)
  214. await resources.dispose()
  215. }
  216. })
  217. it('retains a late initialization failure after its only caller canceled before waiting', async () => {
  218. const { ctx, owner } = await fixture()
  219. const a = await owner('a')
  220. const b = await owner('b')
  221. const ready = Promise.withResolvers<never>()
  222. const open = vi.fn().mockImplementationOnce(() => ready.promise).mockResolvedValue({ value: 1, close: async () => {} })
  223. const resources = new SessionResources<number>(ctx, { label: 'test', exclusive: true, open })
  224. const controller = new AbortController()
  225. const running = resources.run(a.agent, controller.signal, async value => value)
  226. controller.abort(new Error('cancel before waiting'))
  227. await expect(running).rejects.toThrow('cancel before waiting')
  228. ready.reject(new Error('late browser startup failed'))
  229. await vi.waitFor(() => { expect(resources.available(b.agent)).toBe(true) })
  230. expect(await resources.get(b.agent)).toBe(1)
  231. await resources.dispose()
  232. })
  233. it('reports non-Error cancellation and acquisition failures through cancellable waits', async () => {
  234. const { ctx, owner } = await fixture()
  235. const a = await owner('a')
  236. const entered = Promise.withResolvers<undefined>()
  237. const ready = Promise.withResolvers<undefined>()
  238. const open = vi.fn().mockImplementationOnce(async () => {
  239. entered.resolve(undefined)
  240. await ready.promise
  241. return { value: 1, close: async () => {} }
  242. })
  243. const resources = new SessionResources<number>(ctx, { label: 'test', exclusive: false, open })
  244. const controller = new AbortController()
  245. const waiting = resources.get(a.agent, controller.signal)
  246. const canceled = expect(waiting).rejects.toMatchObject({ message: 'browser operation canceled', cause: 'stop' })
  247. await entered.promise
  248. controller.abort('stop')
  249. await canceled
  250. ready.resolve(undefined)
  251. await resources.dispose()
  252. const b = await owner('b')
  253. const failed = new SessionResources<number>(ctx, { label: 'test', exclusive: false, open: vi.fn().mockRejectedValue('failed to connect') })
  254. await expect(failed.get(b.agent, new AbortController().signal)).rejects.toThrow('failed to connect')
  255. await failed.dispose()
  256. })
  257. it('closes a late acquisition and waits for its shutdown during racing disposals', async () => {
  258. const { ctx, owner } = await fixture()
  259. const a = await owner('a')
  260. const entered = Promise.withResolvers<undefined>()
  261. const release = Promise.withResolvers<undefined>()
  262. const closeEntered = Promise.withResolvers<undefined>()
  263. const closeReleased = Promise.withResolvers<undefined>()
  264. const close = vi.fn(async () => { closeEntered.resolve(undefined); await closeReleased.promise })
  265. const resources = new SessionResources(ctx, {
  266. label: 'test', exclusive: false,
  267. async open() { entered.resolve(undefined); await release.promise; return { value: {}, close } },
  268. })
  269. const acquiring = resources.get(a.agent)
  270. const rejected = expect(acquiring).rejects.toThrow('closing')
  271. await entered.promise
  272. const ownerDisposal = a.dispose()
  273. const providerDisposal = resources.dispose()
  274. expect(resources.dispose()).toBe(providerDisposal)
  275. let disposed = false
  276. void providerDisposal.then(() => { disposed = true })
  277. release.resolve(undefined)
  278. await closeEntered.promise
  279. expect(disposed).toBe(false)
  280. closeReleased.resolve(undefined)
  281. await Promise.all([ownerDisposal, providerDisposal, rejected])
  282. expect(close).toHaveBeenCalledTimes(1)
  283. })
  284. it('interrupts resources before awaiting an operation that needs close to settle', async () => {
  285. const { ctx, owner } = await fixture()
  286. const a = await owner('a')
  287. const running = Promise.withResolvers<undefined>()
  288. const stopped = Promise.withResolvers<undefined>()
  289. const resources = new SessionResources(ctx, {
  290. label: 'test', exclusive: false,
  291. open: async () => ({ value: {}, close: async () => { stopped.resolve(undefined) } }),
  292. })
  293. const call = resources.run(a.agent, new AbortController().signal, async (_resource, signal) => {
  294. running.resolve(undefined)
  295. await stopped.promise
  296. signal.throwIfAborted()
  297. })
  298. const rejected = expect(call).rejects.toThrow('closing')
  299. await running.promise
  300. await resources.dispose()
  301. await rejected
  302. })
  303. it('retains exclusive ownership when resource shutdown fails', async () => {
  304. const { ctx, owner } = await fixture()
  305. const a = await owner('a')
  306. const b = await owner('b')
  307. const resources = new SessionResources(ctx, {
  308. label: 'test', exclusive: true,
  309. open: async () => ({ value: {}, close: async () => { throw new Error('close failed') } }),
  310. })
  311. await resources.get(a.agent)
  312. await a.dispose()
  313. await expect(resources.get(b.agent)).rejects.toThrow('already reserved')
  314. await expect(resources.dispose()).rejects.toThrow('browser cleanup failed')
  315. await expect(resources.get(b.agent)).rejects.toThrow('not a live browser owner')
  316. })
  317. })
  318. it('reports early disposal cleanup failure while retaining the owned resource', async () => {
  319. const { ctx, owner } = await fixture()
  320. const a = await owner('early-close-failure')
  321. const entered = Promise.withResolvers<undefined>()
  322. const stopped = Promise.withResolvers<undefined>()
  323. const warning = vi.spyOn(ctx.logger, 'warn')
  324. const resources = new SessionResources(ctx, {
  325. label: 'early-close', exclusive: true,
  326. open: async () => ({ value: {}, async close() { stopped.resolve(undefined); throw new Error('Close failed') } }),
  327. })
  328. const controller = new AbortController()
  329. const running = resources.run(a.agent, controller.signal, async () => { entered.resolve(undefined); await stopped.promise })
  330. const canceled = expect(running).rejects.toMatchObject({ kind: 'disposed' })
  331. await entered.promise
  332. controller.abort({ kind: 'disposed' })
  333. await canceled
  334. await a.dispose()
  335. expect(warning).toHaveBeenCalledWith(expect.stringContaining('cleanup during Session cancellation failed'))
  336. await expect(resources.dispose()).rejects.toThrow('browser cleanup failed')
  337. warning.mockRestore()
  338. })