api-proxy-workspace.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325
  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. acceptsNextStep: false,
  42. ctx: new Context(),
  43. followup: () => 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. pickDirectory?: (signal: AbortSignal) => Promise<string | null>,
  55. ) {
  56. const ctx = new Context()
  57. await ctx.plugin(SessionStore)
  58. await ctx.plugin(AgentRegistry)
  59. await ctx.plugin(UserInteractionService)
  60. await ctx.plugin(Storage)
  61. ctx.storage.backend.register('memory', new MemoryStorageBackend())
  62. const storageDomain = new DomainFacility(ctx, { backend: 'memory', routes: {} })
  63. ctx.storage.mount('domain', storageDomain)
  64. ctx.provide('storageDomain', storageDomain)
  65. ctx.provide('sessionPersistence', { list: () => Promise.resolve([]) } as never)
  66. await ctx.plugin(WorkspaceRegistry)
  67. const factory: AgentFactory = {
  68. async createAgent(_ownerCtx, options) {
  69. const session = ctx.sessions.create(
  70. options.sessionId,
  71. options.meta === undefined ? {} : { meta: options.meta },
  72. )
  73. const agent = stubAgent(session)
  74. const unregister = ctx.agents.register(agent)
  75. return {
  76. agent,
  77. dispose: () => {
  78. unregister()
  79. return Promise.resolve()
  80. },
  81. }
  82. },
  83. async resume() {
  84. throw new Error('test harness has no persisted sessions')
  85. },
  86. }
  87. ctx.agents.setFactory(factory)
  88. const api = createApiProxy(ctx, {
  89. provider: 'test',
  90. model: 'test-model',
  91. cwd: workspaceRoot,
  92. workspaceRoot,
  93. ...pickDirectory === undefined ? {} : { pickDirectory },
  94. })
  95. return { api, ctx, storageDomain, workspaceRoot }
  96. }
  97. describe('host.pickDirectory', () => {
  98. it('returns a selected path or explicit cancellation from the injected native boundary', async () => {
  99. const selected = await harness(undefined, async () => '/tmp/project')
  100. expect((await selected.api.host.pickDirectory(request({}), new AbortController().signal)).result)
  101. .toEqual({ ok: true, value: { path: '/tmp/project' } })
  102. const cancelled = await harness(undefined, async () => null)
  103. expect((await cancelled.api.host.pickDirectory(request({}), new AbortController().signal)).result)
  104. .toEqual({ ok: true, value: { path: null } })
  105. })
  106. it('propagates abort into the native boundary as a cancelled RPC error', async () => {
  107. const { api } = await harness(undefined, signal => new Promise((_resolve, reject) => {
  108. signal.addEventListener('abort', () => { reject(new Error('aborted')) }, { once: true })
  109. }))
  110. const abort = new AbortController()
  111. const pending = api.host.pickDirectory(request({}), abort.signal)
  112. abort.abort()
  113. expect((await pending).result).toMatchObject({ ok: false, error: { code: 'cancelled' } })
  114. })
  115. })
  116. describe('workspace.create', () => {
  117. it('serializes concurrent names and rejects the duplicate', async () => {
  118. const { api, workspaceRoot } = await harness()
  119. const responses = await Promise.all([
  120. api.workspace.create(request({ name: 'alpha' })),
  121. api.workspace.create(request({ name: 'alpha' })),
  122. ])
  123. const created = responses.find(response => response.result.ok)
  124. const duplicate = responses.find(response => !response.result.ok)
  125. expect(created).toBeDefined()
  126. expect(expectOk(created!)).toMatchObject({
  127. created: true,
  128. workspace: { path: join(workspaceRoot, 'alpha'), title: 'alpha' },
  129. })
  130. expect(duplicate?.result).toMatchObject({
  131. ok: false,
  132. error: { code: 'workspace-name-conflict', details: { name: 'alpha' } },
  133. })
  134. expect(existsSync(join(workspaceRoot, 'alpha'))).toBe(true)
  135. })
  136. it('adopts only existing directories and rejects unsafe names', async () => {
  137. const { api, workspaceRoot } = await harness()
  138. const existing = join(workspaceRoot, 'existing')
  139. mkdirSync(existing)
  140. const first = expectOk(await api.workspace.create(request({ path: existing })))
  141. const repeated = expectOk(await api.workspace.create(request({ path: existing })))
  142. expect(first).toMatchObject({ created: true, workspace: { path: existing, title: 'existing' } })
  143. expect(repeated).toMatchObject({ created: false, workspace: { workspaceId: first.workspace.workspaceId } })
  144. expectOk(await api.workspace.rename(request({
  145. workspaceId: first.workspace.workspaceId,
  146. title: 'renamed-existing',
  147. })))
  148. const reopened = expectOk(await api.workspace.create(request({ path: existing })))
  149. expect(reopened.workspace.title).toBe('renamed-existing')
  150. const missing = join(workspaceRoot, 'missing')
  151. const missingResult = await api.workspace.create(request({ path: missing }))
  152. expect(missingResult.result).toMatchObject({ ok: false, error: { code: 'workspace-invalid-path' } })
  153. expect(existsSync(missing)).toBe(false)
  154. for (const name of ['', '.', '..', 'a/b', 'a\\b']) {
  155. const invalid = await api.workspace.create(request({ name }))
  156. expect(invalid.result).toMatchObject({ ok: false, error: { code: 'workspace-invalid-path' } })
  157. }
  158. })
  159. it('rejects different paths that derive the same Workspace title', async () => {
  160. const { api, workspaceRoot } = await harness()
  161. const first = join(workspaceRoot, 'one', 'project')
  162. const second = join(workspaceRoot, 'two', 'project')
  163. mkdirSync(first, { recursive: true })
  164. mkdirSync(second, { recursive: true })
  165. expectOk(await api.workspace.create(request({ path: first })))
  166. const conflict = await api.workspace.create(request({ path: second }))
  167. expect(conflict.result).toMatchObject({
  168. ok: false,
  169. error: { code: 'workspace-name-conflict', details: { name: 'project' } },
  170. })
  171. })
  172. })
  173. describe('session creation and Workspace membership', () => {
  174. it('attaches a preallocated idempotent session while cwd-only sessions stay ungrouped', async () => {
  175. const { api, ctx } = await harness()
  176. const workspace = expectOk(await api.workspace.create(request({ name: 'project' }))).workspace
  177. const sessionId = SessionId('session-workspace-preallocated')
  178. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  179. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  180. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  181. expect(ctx.agents.list().filter(agent => agent.id === sessionId)).toHaveLength(1)
  182. const ungrouped = SessionId('session-cwd-only')
  183. expectOk(await api.sessions.create(request({ cwd: workspace.path, sessionId: ungrouped })))
  184. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  185. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(ungrouped)
  186. const conflict = await api.sessions.create(request({ cwd: join(workspace.path, 'other'), sessionId }))
  187. expect(conflict.result).toMatchObject({
  188. ok: false,
  189. error: { code: 'session-conflict', details: { sessionId, existingCwd: workspace.path } },
  190. })
  191. const missing = await api.sessions.create(request({
  192. workspaceId: 'missing-workspace' as WorkspaceId,
  193. sessionId: SessionId('session-missing-workspace'),
  194. }))
  195. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace-not-found' } })
  196. })
  197. it('retains a published session when attachment fails and repairs it on retry', async () => {
  198. const { api, ctx } = await harness()
  199. const created = expectOk(await api.workspace.create(request({ name: 'project' }))).workspace
  200. const workspace = ctx.workspace.list()[0]
  201. if (workspace === undefined) throw new Error('workspace missing from registry')
  202. vi.spyOn(workspace, 'attachSession').mockRejectedValueOnce(new Error('simulated write failure'))
  203. const sessionId = SessionId('session-attach-retry')
  204. const failed = await api.sessions.create(request({ workspaceId: created.workspaceId, sessionId }))
  205. expect(failed.result).toMatchObject({
  206. ok: false,
  207. error: { code: 'workspace-attach-failed', details: { sessionId, workspaceId: created.workspaceId } },
  208. })
  209. expect(ctx.agents.get(sessionId)).toBeDefined()
  210. expectOk(await api.sessions.create(request({ workspaceId: created.workspaceId, sessionId })))
  211. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  212. })
  213. })
  214. describe('Host Workspace increments', () => {
  215. it('streams committed Workspace and Session increments after empty baselines', async () => {
  216. const { api } = await harness()
  217. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  218. expect(expectOk(await api.sessions.list(request({}))).items).toEqual([])
  219. const abort = new AbortController()
  220. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  221. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  222. const workspaceIncrement = nextHostFrame(stream)
  223. const workspace = expectOk(await api.workspace.create(request({ name: 'project' }))).workspace
  224. expect(await workspaceIncrement).toMatchObject({
  225. payload: { type: 'host/workspace-changed', workspace: { workspaceId: workspace.workspaceId } },
  226. })
  227. const sessionId = SessionId('session-streamed-workspace')
  228. const pending = nextHostFrame(stream)
  229. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  230. const increments: HostFrame[] = []
  231. increments.push((await pending).payload)
  232. while (increments.length < 2) {
  233. const next = await stream.next()
  234. if (next.done === true) throw new Error('Host stream ended before both increments')
  235. increments.push(next.value.payload)
  236. }
  237. expect(increments.find(increment => increment.type === 'host/session-added')).toMatchObject({
  238. // A just-created session has no events: the frame constantly carries blank:true.
  239. type: 'host/session-added', sessionId, blank: true, cwd: workspace.path,
  240. })
  241. const workspaceChanged = increments.find(
  242. (increment): increment is Extract<HostFrame, { type: 'host/workspace-changed' }> =>
  243. increment.type === 'host/workspace-changed',
  244. )
  245. expect(workspaceChanged?.workspace.sessionIds).toEqual([sessionId])
  246. abort.abort()
  247. })
  248. it('does not publish a Workspace whose registry-order commit fails', async () => {
  249. const { api, storageDomain } = await harness()
  250. const domain = storageDomain.get('workspace')
  251. if (domain === undefined) throw new Error('workspace domain is not open')
  252. vi.spyOn(domain.global, 'set').mockRejectedValueOnce(new Error('simulated registry order failure'))
  253. const abort = new AbortController()
  254. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  255. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  256. const next = stream.next()
  257. const failed = await api.workspace.create(request({ name: 'ghost' }))
  258. expect(failed.result.ok).toBe(false)
  259. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  260. abort.abort()
  261. expect(await next).toMatchObject({ done: true })
  262. })
  263. it('deletes the registration, keeps its session and folder, and streams one removal', async () => {
  264. const { api, ctx } = await harness()
  265. const workspace = expectOk(await api.workspace.create(request({ name: 'delete-me' }))).workspace
  266. const sessionId = SessionId('session-kept-after-workspace-delete')
  267. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  268. const abort = new AbortController()
  269. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  270. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  271. const removed = nextHostFrame(stream)
  272. expectOk(await api.workspace.delete(request({ workspaceId: workspace.workspaceId })))
  273. expect(await removed).toMatchObject({
  274. payload: { type: 'host/workspace-removed', workspaceId: workspace.workspaceId },
  275. })
  276. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  277. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(sessionId)
  278. expect(ctx.agents.get(sessionId)).toBeDefined()
  279. expect(existsSync(workspace.path)).toBe(true)
  280. const missing = await api.workspace.delete(request({ workspaceId: workspace.workspaceId }))
  281. expect(missing.result).toMatchObject({
  282. ok: false,
  283. error: { code: 'workspace-not-found', details: { workspaceId: workspace.workspaceId } },
  284. })
  285. const reregistered = expectOk(await api.workspace.create(request({ path: workspace.path }))).workspace
  286. expect(reregistered.workspaceId).not.toBe(workspace.workspaceId)
  287. expect(reregistered.path).toBe(workspace.path)
  288. expect(reregistered.sessionIds).toEqual([])
  289. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(sessionId)
  290. abort.abort()
  291. })
  292. })