api-proxy-workspace.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246
  1. import { existsSync, mkdirSync, mkdtempSync, realpathSync } from 'node:fs'
  2. import { tmpdir } from 'node:os'
  3. import { join } from 'node:path'
  4. import { describe, expect, it, vi } from 'vitest'
  5. import { Context } from 'cordis'
  6. import AgentRegistry, { AgentMessageId } from '@deepseek-ai/dsh-agent'
  7. import type { Agent, AgentFactory } from '@deepseek-ai/dsh-agent'
  8. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  9. import type { Session } from '@deepseek-ai/dsh-session'
  10. import Storage from '@deepseek-ai/dsh-storage'
  11. import { DomainFacility } from '@deepseek-ai/dsh-storage-domain'
  12. import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
  13. import WorkspaceRegistry from '@deepseek-ai/dsh-workspace'
  14. import type { HostFrame, WorkspaceId } from '@deepseek-ai/dsh-host-apiproxy/api'
  15. import type { RpcRequest, RpcResponse } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  16. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  17. import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
  18. import { MemoryStorageBackend } from '../../../storage/storage-domain/tests/helpers/memory-backend.ts'
  19. let nextRpc = 1
  20. function request<P>(payload: P): RpcRequest<P> {
  21. return { rpcId: RpcId(`workspace-${String(nextRpc++)}`), payload }
  22. }
  23. function expectOk<T>(response: RpcResponse<T>): T {
  24. expect(response.result.ok).toBe(true)
  25. if (!response.result.ok) throw new Error('unreachable')
  26. return response.result.value
  27. }
  28. async function nextHostFrame(
  29. stream: AsyncIterator<RpcRequest<HostFrame>>,
  30. ): Promise<RpcRequest<HostFrame>> {
  31. const next = await stream.next()
  32. if (next.done === true) throw new Error('Host stream ended before the expected increment')
  33. return next.value
  34. }
  35. function stubAgent(session: Session): Agent {
  36. return {
  37. id: session.id,
  38. options: {},
  39. session,
  40. status: 'idle',
  41. ctx: new Context(),
  42. followup: () => AgentMessageId('stub'),
  43. queue: () => AgentMessageId('stub'),
  44. steer: () => AgentMessageId('stub'),
  45. inject: () => AgentMessageId('stub'),
  46. send: () => AgentMessageId('stub'),
  47. cancel() {},
  48. whenIdle: () => Promise.resolve(),
  49. }
  50. }
  51. /** Compose the API over real Session, Agent, Storage, Domain, and Workspace services. */
  52. async function harness(
  53. workspaceRoot = realpathSync(mkdtempSync(join(tmpdir(), 'dsh-apiproxy-workspace-'))),
  54. ) {
  55. const ctx = new Context()
  56. await ctx.plugin(SessionStore)
  57. await ctx.plugin(AgentRegistry)
  58. await ctx.plugin(UserInteractionService)
  59. await ctx.plugin(Storage)
  60. ctx.storage.backend.register('memory', new MemoryStorageBackend())
  61. const storageDomain = new DomainFacility(ctx, { backend: 'memory', routes: {} })
  62. ctx.storage.mount('domain', storageDomain)
  63. ctx.provide('storageDomain', storageDomain)
  64. ctx.provide('sessionPersistence', { list: () => Promise.resolve([]) } as never)
  65. await ctx.plugin(WorkspaceRegistry)
  66. const factory: AgentFactory = {
  67. async createAgent(_ownerCtx, options) {
  68. const session = ctx.sessions.create(
  69. options.sessionId,
  70. options.meta === undefined ? {} : { meta: options.meta },
  71. )
  72. const agent = stubAgent(session)
  73. const unregister = ctx.agents.register(agent)
  74. return {
  75. agent,
  76. dispose: () => {
  77. unregister()
  78. return Promise.resolve()
  79. },
  80. }
  81. },
  82. async resume() {
  83. throw new Error('test harness has no persisted sessions')
  84. },
  85. }
  86. ctx.agents.setFactory(factory)
  87. const api = createApiProxy(ctx, {
  88. provider: 'test',
  89. model: 'test-model',
  90. cwd: workspaceRoot,
  91. workspaceRoot,
  92. })
  93. return { api, ctx, storageDomain, workspaceRoot }
  94. }
  95. describe('workspace.create', () => {
  96. it('serializes concurrent names and rejects the duplicate', async () => {
  97. const { api, workspaceRoot } = await harness()
  98. const responses = await Promise.all([
  99. api.workspace.create(request({ name: 'alpha' })),
  100. api.workspace.create(request({ name: 'alpha' })),
  101. ])
  102. const created = responses.find(response => response.result.ok)
  103. const duplicate = responses.find(response => !response.result.ok)
  104. expect(created).toBeDefined()
  105. expect(expectOk(created!)).toMatchObject({
  106. created: true,
  107. workspace: { path: join(workspaceRoot, 'alpha'), title: 'alpha' },
  108. })
  109. expect(duplicate?.result).toMatchObject({
  110. ok: false,
  111. error: { code: 'workspace-name-conflict', details: { name: 'alpha' } },
  112. })
  113. expect(existsSync(join(workspaceRoot, 'alpha'))).toBe(true)
  114. })
  115. it('adopts only existing directories and rejects unsafe names', async () => {
  116. const { api, workspaceRoot } = await harness()
  117. const existing = join(workspaceRoot, 'existing')
  118. mkdirSync(existing)
  119. const first = expectOk(await api.workspace.create(request({ path: existing })))
  120. const repeated = expectOk(await api.workspace.create(request({ path: existing })))
  121. expect(first).toMatchObject({ created: true, workspace: { path: existing, title: 'existing' } })
  122. expect(repeated).toMatchObject({ created: false, workspace: { workspaceId: first.workspace.workspaceId } })
  123. const missing = join(workspaceRoot, 'missing')
  124. const missingResult = await api.workspace.create(request({ path: missing }))
  125. expect(missingResult.result).toMatchObject({ ok: false, error: { code: 'workspace-invalid-path' } })
  126. expect(existsSync(missing)).toBe(false)
  127. for (const name of ['', '.', '..', 'a/b', 'a\\b']) {
  128. const invalid = await api.workspace.create(request({ name }))
  129. expect(invalid.result).toMatchObject({ ok: false, error: { code: 'workspace-invalid-path' } })
  130. }
  131. })
  132. })
  133. describe('session creation and Workspace membership', () => {
  134. it('attaches a preallocated idempotent session while cwd-only sessions stay ungrouped', async () => {
  135. const { api, ctx } = await harness()
  136. const workspace = expectOk(await api.workspace.create(request({ name: 'project' }))).workspace
  137. const sessionId = SessionId('session-workspace-preallocated')
  138. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  139. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  140. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  141. expect(ctx.agents.list().filter(agent => agent.id === sessionId)).toHaveLength(1)
  142. const ungrouped = SessionId('session-cwd-only')
  143. expectOk(await api.sessions.create(request({ cwd: workspace.path, sessionId: ungrouped })))
  144. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  145. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(ungrouped)
  146. const conflict = await api.sessions.create(request({ cwd: join(workspace.path, 'other'), sessionId }))
  147. expect(conflict.result).toMatchObject({
  148. ok: false,
  149. error: { code: 'session-conflict', details: { sessionId, existingCwd: workspace.path } },
  150. })
  151. const missing = await api.sessions.create(request({
  152. workspaceId: 'missing-workspace' as WorkspaceId,
  153. sessionId: SessionId('session-missing-workspace'),
  154. }))
  155. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace-not-found' } })
  156. })
  157. it('retains a published session when attachment fails and repairs it on retry', async () => {
  158. const { api, ctx } = await harness()
  159. const created = expectOk(await api.workspace.create(request({ name: 'project' }))).workspace
  160. const workspace = ctx.workspace.list()[0]
  161. if (workspace === undefined) throw new Error('workspace missing from registry')
  162. vi.spyOn(workspace, 'attachSession').mockRejectedValueOnce(new Error('simulated write failure'))
  163. const sessionId = SessionId('session-attach-retry')
  164. const failed = await api.sessions.create(request({ workspaceId: created.workspaceId, sessionId }))
  165. expect(failed.result).toMatchObject({
  166. ok: false,
  167. error: { code: 'workspace-attach-failed', details: { sessionId, workspaceId: created.workspaceId } },
  168. })
  169. expect(ctx.agents.get(sessionId)).toBeDefined()
  170. expectOk(await api.sessions.create(request({ workspaceId: created.workspaceId, sessionId })))
  171. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  172. })
  173. })
  174. describe('Host Workspace increments', () => {
  175. it('streams committed Workspace and Session increments after empty baselines', async () => {
  176. const { api } = await harness()
  177. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  178. expect(expectOk(await api.sessions.list(request({}))).items).toEqual([])
  179. const abort = new AbortController()
  180. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  181. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  182. const workspaceIncrement = nextHostFrame(stream)
  183. const workspace = expectOk(await api.workspace.create(request({ name: 'project' }))).workspace
  184. expect(await workspaceIncrement).toMatchObject({
  185. payload: { type: 'host/workspace-changed', workspace: { workspaceId: workspace.workspaceId } },
  186. })
  187. const sessionId = SessionId('session-streamed-workspace')
  188. const pending = nextHostFrame(stream)
  189. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  190. const increments: HostFrame[] = []
  191. increments.push((await pending).payload)
  192. while (increments.length < 2) {
  193. const next = await stream.next()
  194. if (next.done === true) throw new Error('Host stream ended before both increments')
  195. increments.push(next.value.payload)
  196. }
  197. expect(increments.find(increment => increment.type === 'host/session-added')).toMatchObject({
  198. type: 'host/session-added', sessionId, cwd: workspace.path,
  199. })
  200. const workspaceChanged = increments.find(
  201. (increment): increment is Extract<HostFrame, { type: 'host/workspace-changed' }> =>
  202. increment.type === 'host/workspace-changed',
  203. )
  204. expect(workspaceChanged?.workspace.sessionIds).toEqual([sessionId])
  205. abort.abort()
  206. })
  207. it('does not publish a Workspace whose registry-order commit fails', async () => {
  208. const { api, storageDomain } = await harness()
  209. const domain = storageDomain.get('workspace')
  210. if (domain === undefined) throw new Error('workspace domain is not open')
  211. vi.spyOn(domain.global, 'set').mockRejectedValueOnce(new Error('simulated registry order failure'))
  212. const abort = new AbortController()
  213. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  214. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  215. const next = stream.next()
  216. const failed = await api.workspace.create(request({ name: 'ghost' }))
  217. expect(failed.result.ok).toBe(false)
  218. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  219. abort.abort()
  220. expect(await next).toMatchObject({ done: true })
  221. })
  222. })