api-proxy-workspace.spec.ts 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636
  1. import { existsSync, mkdirSync, mkdtempSync, realpathSync } from 'node:fs'
  2. import { homedir, tmpdir } from 'node:os'
  3. import { join } from 'node:path'
  4. import { describe, expect, it, vi } from 'vitest'
  5. import { Context } from '@deepseek-ai/cordis'
  6. import AgentRegistry, { agentEvents, Inbox } 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 UserQuestionService from '@deepseek-ai/dsh-user-questions'
  13. import { DirectoryPickerError } from '@deepseek-ai/dsh-host-directory-picker'
  14. import type { DirectoryPickerCapability } from '@deepseek-ai/dsh-host-directory-picker'
  15. import WorkspaceRegistry from '@deepseek-ai/dsh-workspace'
  16. import type { HostFrame, WorkspaceId } from '@deepseek-ai/dsh-host-apiproxy/api'
  17. import type { RpcRequest, RpcResponse } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  18. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  19. import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
  20. import { MemoryStorageBackend } from '../../../storage/storage-domain/tests/helpers/memory-backend.ts'
  21. let nextRpc = 1
  22. function request<P>(payload: P): RpcRequest<P> {
  23. return { rpcId: RpcId(`workspace-${String(nextRpc++)}`), payload }
  24. }
  25. function expectOk<T>(response: RpcResponse<T>): T {
  26. expect(response.result.ok).toBe(true)
  27. if (!response.result.ok) throw new Error('unreachable')
  28. return response.result.value
  29. }
  30. async function nextHostFrame(
  31. stream: AsyncIterator<RpcRequest<HostFrame>>,
  32. ): Promise<RpcRequest<HostFrame>> {
  33. const next = await stream.next()
  34. if (next.done === true) throw new Error('Host stream ended before the expected increment')
  35. return next.value
  36. }
  37. function stubAgent(session: Session): Agent {
  38. const agent: Agent = {
  39. id: session.id,
  40. options: {},
  41. session,
  42. inbox: undefined as never,
  43. status: 'idle',
  44. ctx: new Context(),
  45. send: () => {},
  46. followup: () => {},
  47. steer: () => ({ outcome: Promise.resolve({ status: 'rejected' as const }) }),
  48. inject: () => {},
  49. cancel() {},
  50. runMaintenance: job => job(new AbortController().signal),
  51. whenIdle: () => Promise.resolve(),
  52. }
  53. Object.assign(agent, { inbox: new Inbox(agent.ctx, agent.session, agentEvents(agent.ctx, agent)) })
  54. return agent
  55. }
  56. /** Compose the API over real Session, Agent, Storage, Domain, and Workspace services. */
  57. async function harness(
  58. root = realpathSync.native(mkdtempSync(join(tmpdir(), 'dsh-apiproxy-workspace-'))),
  59. picker: DirectoryPickerCapability = { kind: 'native', pick: async () => null },
  60. extras: {
  61. openPath?: (path: string, signal: AbortSignal) => Promise<void>
  62. canOpenPath?: () => boolean
  63. refreshDefaultForReuse?: (session: Session) => void
  64. } = {},
  65. ) {
  66. const ctx = new Context()
  67. await ctx.plugin(SessionStore)
  68. await ctx.plugin(AgentRegistry)
  69. await ctx.plugin(UserQuestionService)
  70. await ctx.plugin(Storage)
  71. ctx.storage.backend.register('memory', new MemoryStorageBackend())
  72. const storageDomain = new DomainFacility(ctx, { backend: 'memory', routes: {} })
  73. ctx.storage.mount('domain', storageDomain)
  74. ctx.provide('storageDomain', storageDomain)
  75. ctx.provide('sessionPersistence', { list: () => Promise.resolve([]) } as never)
  76. await ctx.plugin(WorkspaceRegistry)
  77. const factory: AgentFactory = {
  78. async createAgent(_ownerCtx, options) {
  79. const session = ctx.sessions.create(
  80. options.sessionId,
  81. options.meta === undefined ? {} : { meta: options.meta },
  82. )
  83. const agent = stubAgent(session)
  84. const unregister = ctx.agents.register(agent)
  85. return {
  86. agent,
  87. dispose: () => {
  88. unregister()
  89. return Promise.resolve()
  90. },
  91. }
  92. },
  93. async resume() {
  94. throw new Error('test harness has no persisted sessions')
  95. },
  96. }
  97. ctx.agents.setFactory(factory)
  98. // Structural picker fake: the gateway only reads capability(); a stable
  99. // object per harness mirrors the seam's stability contract.
  100. ctx.provide('directoryPicker', { capability: () => picker } as never)
  101. if (extras.refreshDefaultForReuse !== undefined) {
  102. ctx.provide('permissionPresets', {
  103. refreshDefaultForReuse: extras.refreshDefaultForReuse,
  104. } as never)
  105. }
  106. const api = createApiProxy(ctx, {
  107. defaultModelSelection: () => ({ provider: 'test', model: 'test-model' }),
  108. cwd: root,
  109. ...extras.openPath === undefined ? {} : { openPath: extras.openPath },
  110. ...extras.canOpenPath === undefined ? {} : { canOpenPath: extras.canOpenPath },
  111. })
  112. return { api, ctx, storageDomain, root }
  113. }
  114. /** Stage one directory under the harness root for path adoption. */
  115. function stageDir(root: string, name: string): string {
  116. const path = join(root, name)
  117. mkdirSync(path)
  118. return path
  119. }
  120. describe('host.pickDirectory', () => {
  121. it('returns a selected path or explicit cancellation from the native capability', async () => {
  122. const selected = await harness(undefined, { kind: 'native', pick: async () => '/tmp/project' })
  123. expect((await selected.api.host.pickDirectory(request({}), new AbortController().signal)).result)
  124. .toEqual({ ok: true, value: { path: '/tmp/project' } })
  125. const cancelled = await harness(undefined, { kind: 'native', pick: async () => null })
  126. expect((await cancelled.api.host.pickDirectory(request({}), new AbortController().signal)).result)
  127. .toEqual({ ok: true, value: { path: null } })
  128. })
  129. it('propagates abort into the native capability as a cancelled RPC error', async () => {
  130. const { api } = await harness(undefined, {
  131. kind: 'native',
  132. pick: signal => new Promise((_resolve, reject) => {
  133. signal.addEventListener('abort', () => { reject(new Error('aborted')) }, { once: true })
  134. }),
  135. })
  136. const abort = new AbortController()
  137. const pending = api.host.pickDirectory(request({}), abort.signal)
  138. abort.abort()
  139. expect((await pending).result).toMatchObject({ ok: false, error: { code: 'cancelled' } })
  140. })
  141. it('folds a non-abort native-chooser failure into an internal error', async () => {
  142. const { api } = await harness(undefined, { kind: 'native', pick: async () => { throw new Error('no chooser installed') } })
  143. const response = await api.host.pickDirectory(request({}), new AbortController().signal)
  144. expect(response.result).toMatchObject({ ok: false, error: { code: 'internal' } })
  145. })
  146. it('refuses the native RPC under a browse composition', async () => {
  147. const { api } = await harness(undefined, BROWSE_STUB)
  148. const response = await api.host.pickDirectory(request({}), new AbortController().signal)
  149. expect(response.result).toMatchObject({
  150. ok: false,
  151. error: { code: 'directory-picker-unavailable', details: { capability: 'browse' } },
  152. })
  153. })
  154. })
  155. /** Canned browse capability: one listing, one created path, typed failures on demand. */
  156. const BROWSE_STUB: DirectoryPickerCapability = {
  157. kind: 'browse',
  158. list: async (path) => {
  159. if (path === '/denied') throw new DirectoryPickerError('directory-unreadable', '/denied', 'cannot list /denied')
  160. const target = path ?? '/home/user'
  161. return {
  162. path: target,
  163. home: '/home/user',
  164. crumbs: [{ name: '/', path: '/', hidden: false }],
  165. entries: [{ name: 'projects', path: `${target}/projects`, hidden: false }],
  166. truncated: false,
  167. }
  168. },
  169. createDirectory: async (path, name) => {
  170. if (name === 'taken') throw new DirectoryPickerError('directory-exists', `${path}/${name}`, 'already exists')
  171. if (name === 'unwritable') throw new Error('disk detached')
  172. return `${path}/${name}`
  173. },
  174. }
  175. describe('host.listDirectory / host.createDirectory', () => {
  176. it('serves listings and creation through the browse capability, defaulting to home', async () => {
  177. const { api } = await harness(undefined, BROWSE_STUB)
  178. const home = await api.host.listDirectory(request({}), new AbortController().signal)
  179. expect(home.result).toMatchObject({ ok: true, value: { path: '/home/user', home: '/home/user' } })
  180. const listed = await api.host.listDirectory(request({ path: '/home/user/projects' }), new AbortController().signal)
  181. expect(listed.result).toMatchObject({ ok: true, value: { path: '/home/user/projects' } })
  182. const created = await api.host.createDirectory(request({ path: '/home/user', name: 'fresh' }))
  183. expect(created.result).toEqual({ ok: true, value: { path: '/home/user/fresh' } })
  184. })
  185. it('maps typed picker failures onto the wire error codes and folds unknown throws to internal', async () => {
  186. const { api } = await harness(undefined, BROWSE_STUB)
  187. expect((await api.host.listDirectory(request({ path: '/denied' }), new AbortController().signal)).result).toMatchObject({
  188. ok: false, error: { code: 'directory-unreadable', details: { path: '/denied' } },
  189. })
  190. expect((await api.host.createDirectory(request({ path: '/home/user', name: 'taken' }))).result).toMatchObject({
  191. ok: false, error: { code: 'directory-exists' },
  192. })
  193. expect((await api.host.createDirectory(request({ path: '/home/user', name: 'unwritable' }))).result).toMatchObject({
  194. ok: false, error: { code: 'internal' },
  195. })
  196. })
  197. it('reports an aborted listing as cancelled, like the other signal-following RPCs', async () => {
  198. const { api } = await harness(undefined, {
  199. kind: 'browse',
  200. list: (_path, signal) => new Promise((_resolve, reject) => {
  201. signal?.addEventListener('abort', () => { reject(new Error('scan aborted')) }, { once: true })
  202. }),
  203. createDirectory: async () => '/never',
  204. })
  205. const abort = new AbortController()
  206. const pending = api.host.listDirectory(request({}), abort.signal)
  207. abort.abort()
  208. expect((await pending).result).toMatchObject({ ok: false, error: { code: 'cancelled' } })
  209. })
  210. it('refuses the browse RPCs under a native composition', async () => {
  211. const { api } = await harness()
  212. expect((await api.host.listDirectory(request({}), new AbortController().signal)).result).toMatchObject({
  213. ok: false, error: { code: 'directory-picker-unavailable', details: { capability: 'native' } },
  214. })
  215. expect((await api.host.createDirectory(request({ path: '/x', name: 'y' }))).result).toMatchObject({
  216. ok: false, error: { code: 'directory-picker-unavailable', details: { capability: 'native' } },
  217. })
  218. })
  219. })
  220. describe('host.openPath', () => {
  221. it('describes whether this deployment can reach a user-visible native desktop', async () => {
  222. const visible = await harness(undefined, undefined, { canOpenPath: () => true })
  223. const headless = await harness(undefined, undefined, { canOpenPath: () => false })
  224. expect(expectOk(await visible.api.host.describe(request({}))).canOpenPath).toBe(true)
  225. expect(expectOk(await headless.api.host.describe(request({}))).canOpenPath).toBe(false)
  226. expect(expectOk(await visible.api.host.describe(request({}))).home).toBe(homedir())
  227. })
  228. it('opens through the injected native boundary', async () => {
  229. const opened: string[] = []
  230. const { api } = await harness(undefined, undefined, {
  231. openPath: async (path) => { opened.push(path) },
  232. })
  233. expect((await api.host.openPath(request({ path: '/tmp/a.txt' }), new AbortController().signal)).result)
  234. .toEqual({ ok: true, value: { opened: true } })
  235. expect(opened).toEqual(['/tmp/a.txt'])
  236. })
  237. it('propagates abort into the native boundary as a cancelled RPC error', async () => {
  238. const { api } = await harness(undefined, undefined, {
  239. openPath: (_path, signal) => new Promise((_resolve, reject) => {
  240. signal.addEventListener('abort', () => { reject(new Error('aborted')) }, { once: true })
  241. }),
  242. })
  243. const abort = new AbortController()
  244. const pending = api.host.openPath(request({ path: '/tmp/a.txt' }), abort.signal)
  245. abort.abort()
  246. expect((await pending).result).toMatchObject({ ok: false, error: { code: 'cancelled' } })
  247. })
  248. })
  249. describe('workspace.create', () => {
  250. it('serializes concurrent creates of one path into a single registration', async () => {
  251. const { api, root } = await harness()
  252. const target = stageDir(root, 'alpha')
  253. const responses = await Promise.all([
  254. api.workspace.create(request({ path: target })),
  255. api.workspace.create(request({ path: target })),
  256. ])
  257. const values = responses.map(response => expectOk(response))
  258. const created = values.find(value => value.created)
  259. const resolved = values.find(value => !value.created)
  260. expect(created).toMatchObject({ workspace: { path: target, title: 'alpha' } })
  261. expect(resolved?.workspace.workspaceId).toBe(created?.workspace.workspaceId)
  262. expect(expectOk(await api.workspace.list(request({}))).items).toHaveLength(1)
  263. })
  264. it('adopts only existing directories', async () => {
  265. const { api, root } = await harness()
  266. const existing = stageDir(root, 'existing')
  267. const first = expectOk(await api.workspace.create(request({ path: existing })))
  268. const repeated = expectOk(await api.workspace.create(request({ path: existing })))
  269. expect(first).toMatchObject({ created: true, workspace: { path: existing, title: 'existing' } })
  270. expect(repeated).toMatchObject({ created: false, workspace: { workspaceId: first.workspace.workspaceId } })
  271. expectOk(await api.workspace.rename(request({
  272. workspaceId: first.workspace.workspaceId,
  273. title: 'renamed-existing',
  274. })))
  275. const reopened = expectOk(await api.workspace.create(request({ path: existing })))
  276. expect(reopened.workspace.title).toBe('renamed-existing')
  277. const missing = join(root, 'missing')
  278. const missingResult = await api.workspace.create(request({ path: missing }))
  279. expect(missingResult.result).toMatchObject({ ok: false, error: { code: 'workspace-invalid-path' } })
  280. expect(existsSync(missing)).toBe(false)
  281. })
  282. it('adopts different paths that derive the same Workspace title', async () => {
  283. const { api, root } = await harness()
  284. const first = join(root, 'one', 'project')
  285. const second = join(root, 'two', 'project')
  286. mkdirSync(first, { recursive: true })
  287. mkdirSync(second, { recursive: true })
  288. const firstResult = expectOk(await api.workspace.create(request({ path: first })))
  289. const secondResult = expectOk(await api.workspace.create(request({ path: second })))
  290. expect(firstResult).toMatchObject({
  291. created: true,
  292. workspace: { path: first, title: 'project' },
  293. })
  294. expect(secondResult).toMatchObject({
  295. created: true,
  296. workspace: { path: second, title: 'project' },
  297. })
  298. expect(secondResult.workspace.workspaceId).not.toBe(firstResult.workspace.workspaceId)
  299. expect(expectOk(await api.workspace.list(request({}))).items.map(workspace => workspace.path))
  300. .toEqual([second, first])
  301. })
  302. })
  303. describe('workspace.insertBefore', () => {
  304. it('commits the complete order, streams one order frame, and maps unknown ids', async () => {
  305. const { api, ctx, root } = await harness()
  306. const first = expectOk(await api.workspace.create(request({ path: stageDir(root, 'first') }))).workspace
  307. const second = expectOk(await api.workspace.create(request({ path: stageDir(root, 'second') }))).workspace
  308. const third = expectOk(await api.workspace.create(request({ path: stageDir(root, 'third') }))).workspace
  309. const abort = new AbortController()
  310. const listWorkspaces = vi.spyOn(ctx.workspaceRegistry, 'list')
  311. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  312. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  313. expect(listWorkspaces).toHaveBeenCalledTimes(1)
  314. const changed = nextHostFrame(stream)
  315. const reordered = expectOk(await api.workspace.insertBefore(request({
  316. workspaceId: first.workspaceId,
  317. beforeWorkspaceId: second.workspaceId,
  318. })))
  319. expect(reordered.workspaceIds).toEqual([third.workspaceId, first.workspaceId, second.workspaceId])
  320. expect(await changed).toMatchObject({
  321. payload: {
  322. type: 'host/workspace-order-changed',
  323. workspaceIds: [third.workspaceId, first.workspaceId, second.workspaceId],
  324. },
  325. })
  326. expect(expectOk(await api.workspace.list(request({}))).items.map(item => item.workspaceId))
  327. .toEqual(reordered.workspaceIds)
  328. const missingSource = await api.workspace.insertBefore(request({
  329. workspaceId: 'missing' as WorkspaceId,
  330. }))
  331. expect(missingSource.result).toMatchObject({
  332. ok: false, error: { code: 'workspace-not-found', details: { workspaceId: 'missing' } },
  333. })
  334. const missingAnchor = await api.workspace.insertBefore(request({
  335. workspaceId: first.workspaceId,
  336. beforeWorkspaceId: 'missing-anchor' as WorkspaceId,
  337. }))
  338. expect(missingAnchor.result).toMatchObject({
  339. ok: false, error: { code: 'workspace-not-found', details: { workspaceId: 'missing-anchor' } },
  340. })
  341. abort.abort()
  342. })
  343. })
  344. describe('session creation and Workspace membership', () => {
  345. it('notifies the permission owner only while the confirmed reuse target remains eligible', async () => {
  346. const refreshDefaultForReuse = vi.fn<(session: Session) => void>()
  347. const { api, ctx, root } = await harness(undefined, undefined, { refreshDefaultForReuse })
  348. const workspace = expectOk(await api.workspace.create(request({
  349. path: stageDir(root, 'permission-refresh'),
  350. }))).workspace
  351. const reusedId = SessionId('session-reused-blank')
  352. expectOk(await api.sessions.create(request({
  353. workspaceId: workspace.workspaceId,
  354. sessionId: reusedId,
  355. })))
  356. expect(refreshDefaultForReuse).not.toHaveBeenCalled()
  357. expectOk(await api.sessions.create(request({
  358. workspaceId: workspace.workspaceId,
  359. sessionId: reusedId,
  360. reuseWorkspaceBlank: true,
  361. })))
  362. expect(refreshDefaultForReuse).toHaveBeenCalledOnce()
  363. expect(refreshDefaultForReuse.mock.calls[0]?.[0].id).toBe(reusedId)
  364. const reused = ctx.sessions.get(reusedId)
  365. if (reused === undefined) throw new Error('reused session was not published')
  366. reused.append('turn/start', { turn: 1 })
  367. expectOk(await api.sessions.create(request({
  368. workspaceId: workspace.workspaceId,
  369. sessionId: reusedId,
  370. reuseWorkspaceBlank: true,
  371. })))
  372. expect(refreshDefaultForReuse).toHaveBeenCalledOnce()
  373. const archivedId = SessionId('session-archived-blank')
  374. expectOk(await api.sessions.create(request({
  375. workspaceId: workspace.workspaceId,
  376. sessionId: archivedId,
  377. })))
  378. expectOk(await api.workspace.archiveSession(request({ sessionId: archivedId })))
  379. expectOk(await api.sessions.create(request({
  380. workspaceId: workspace.workspaceId,
  381. sessionId: archivedId,
  382. reuseWorkspaceBlank: true,
  383. })))
  384. expect(refreshDefaultForReuse).toHaveBeenCalledOnce()
  385. const nonMemberId = SessionId('session-non-member-blank')
  386. expectOk(await api.sessions.create(request({
  387. cwd: workspace.path,
  388. sessionId: nonMemberId,
  389. })))
  390. expectOk(await api.sessions.create(request({
  391. workspaceId: workspace.workspaceId,
  392. sessionId: nonMemberId,
  393. reuseWorkspaceBlank: true,
  394. })))
  395. expect(refreshDefaultForReuse).toHaveBeenCalledOnce()
  396. })
  397. it('attaches a preallocated idempotent session while cwd-only sessions stay ungrouped', async () => {
  398. const { api, ctx, root } = await harness()
  399. const workspace = expectOk(await api.workspace.create(request({ path: stageDir(root, 'project') }))).workspace
  400. const sessionId = SessionId('session-workspace-preallocated')
  401. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  402. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  403. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  404. expect(ctx.agents.list().filter(agent => agent.id === sessionId)).toHaveLength(1)
  405. const ungrouped = SessionId('session-cwd-only')
  406. expectOk(await api.sessions.create(request({ cwd: workspace.path, sessionId: ungrouped })))
  407. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  408. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(ungrouped)
  409. const conflict = await api.sessions.create(request({ cwd: join(workspace.path, 'other'), sessionId }))
  410. expect(conflict.result).toMatchObject({
  411. ok: false,
  412. error: { code: 'session-conflict', details: { sessionId, existingCwd: workspace.path } },
  413. })
  414. const missing = await api.sessions.create(request({
  415. workspaceId: 'missing-workspace' as WorkspaceId,
  416. sessionId: SessionId('session-missing-workspace'),
  417. }))
  418. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace-not-found' } })
  419. })
  420. it('retains a published session when attachment fails and repairs it on retry', async () => {
  421. const { api, ctx, root } = await harness()
  422. const created = expectOk(await api.workspace.create(request({ path: stageDir(root, 'project') }))).workspace
  423. const workspace = ctx.workspaceRegistry.list()[0]
  424. if (workspace === undefined) throw new Error('workspace missing from registry')
  425. vi.spyOn(workspace, 'attachSession').mockRejectedValueOnce(new Error('simulated write failure'))
  426. const sessionId = SessionId('session-attach-retry')
  427. const failed = await api.sessions.create(request({ workspaceId: created.workspaceId, sessionId }))
  428. expect(failed.result).toMatchObject({
  429. ok: false,
  430. error: { code: 'workspace-attach-failed', details: { sessionId, workspaceId: created.workspaceId } },
  431. })
  432. expect(ctx.agents.get(sessionId)).toBeDefined()
  433. expectOk(await api.sessions.create(request({ workspaceId: created.workspaceId, sessionId })))
  434. expect(expectOk(await api.workspace.list(request({}))).items[0]?.sessionIds).toEqual([sessionId])
  435. })
  436. })
  437. describe('Host Workspace increments', () => {
  438. it('projects subagent origin in attached summaries and creation increments', async () => {
  439. const { api, ctx } = await harness()
  440. const abort = new AbortController()
  441. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  442. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  443. const pending = nextHostFrame(stream)
  444. const childId = SessionId('session-subagent-child')
  445. ctx.sessions.create(childId, {
  446. meta: {
  447. cwd: '/tmp',
  448. parentSession: SessionId('session-parent'),
  449. origin: 'subagent',
  450. },
  451. })
  452. expect(await pending).toMatchObject({
  453. payload: {
  454. type: 'host/session-added',
  455. sessionId: childId,
  456. parentSessionId: 'session-parent',
  457. origin: 'subagent',
  458. },
  459. })
  460. expect(expectOk(await api.sessions.list(request({}))).items).toContainEqual(
  461. expect.objectContaining({ sessionId: childId, origin: 'subagent' }),
  462. )
  463. abort.abort()
  464. })
  465. it('streams committed Workspace and Session increments after empty baselines', async () => {
  466. const { api, root } = await harness()
  467. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  468. expect(expectOk(await api.sessions.list(request({}))).items).toEqual([])
  469. const abort = new AbortController()
  470. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  471. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  472. const workspaceIncrement = nextHostFrame(stream)
  473. const workspace = expectOk(await api.workspace.create(request({ path: stageDir(root, 'project') }))).workspace
  474. expect(await workspaceIncrement).toMatchObject({
  475. payload: { type: 'host/workspace-changed', workspace: { workspaceId: workspace.workspaceId } },
  476. })
  477. const sessionId = SessionId('session-streamed-workspace')
  478. const pending = nextHostFrame(stream)
  479. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  480. const increments: HostFrame[] = []
  481. increments.push((await pending).payload)
  482. while (increments.length < 2) {
  483. const next = await stream.next()
  484. if (next.done === true) throw new Error('Host stream ended before both increments')
  485. increments.push(next.value.payload)
  486. }
  487. expect(increments.find(increment => increment.type === 'host/session-added')).toMatchObject({
  488. // A just-created session has no events: the frame constantly carries blank:true.
  489. type: 'host/session-added', sessionId, blank: true, cwd: workspace.path,
  490. })
  491. const workspaceChanged = increments.find(
  492. (increment): increment is Extract<HostFrame, { type: 'host/workspace-changed' }> =>
  493. increment.type === 'host/workspace-changed',
  494. )
  495. expect(workspaceChanged?.workspace.sessionIds).toEqual([sessionId])
  496. abort.abort()
  497. })
  498. it('does not publish a Workspace whose registry-order commit fails', async () => {
  499. const { api, storageDomain, root } = await harness()
  500. const domain = storageDomain.get('workspace')
  501. if (domain === undefined) throw new Error('workspace domain is not open')
  502. vi.spyOn(domain.global, 'set').mockRejectedValueOnce(new Error('simulated registry order failure'))
  503. const abort = new AbortController()
  504. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  505. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  506. const next = stream.next()
  507. const failed = await api.workspace.create(request({ path: stageDir(root, 'ghost') }))
  508. expect(failed.result.ok).toBe(false)
  509. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  510. abort.abort()
  511. expect(await next).toMatchObject({ done: true })
  512. })
  513. it('deletes the registration, keeps its session and folder, and streams one removal', async () => {
  514. const { api, ctx, root } = await harness()
  515. const workspace = expectOk(await api.workspace.create(request({ path: stageDir(root, 'delete-me') }))).workspace
  516. const sessionId = SessionId('session-kept-after-workspace-delete')
  517. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  518. const abort = new AbortController()
  519. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  520. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  521. const removed = nextHostFrame(stream)
  522. expectOk(await api.workspace.delete(request({ workspaceId: workspace.workspaceId })))
  523. expect(await removed).toMatchObject({
  524. payload: { type: 'host/workspace-removed', workspaceId: workspace.workspaceId },
  525. })
  526. expect(expectOk(await api.workspace.list(request({}))).items).toEqual([])
  527. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(sessionId)
  528. expect(ctx.agents.get(sessionId)).toBeDefined()
  529. expect(existsSync(workspace.path)).toBe(true)
  530. const missing = await api.workspace.delete(request({ workspaceId: workspace.workspaceId }))
  531. expect(missing.result).toMatchObject({
  532. ok: false,
  533. error: { code: 'workspace-not-found', details: { workspaceId: workspace.workspaceId } },
  534. })
  535. const reregistered = expectOk(await api.workspace.create(request({ path: workspace.path }))).workspace
  536. expect(reregistered.workspaceId).not.toBe(workspace.workspaceId)
  537. expect(reregistered.path).toBe(workspace.path)
  538. expect(reregistered.sessionIds).toEqual([])
  539. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(sessionId)
  540. abort.abort()
  541. })
  542. it('archives a session into the global set, keeps its accounting, and streams the set once', async () => {
  543. const { api, root } = await harness()
  544. const workspace = expectOk(await api.workspace.create(request({ path: stageDir(root, 'archive-home') }))).workspace
  545. const sessionId = SessionId('session-to-archive')
  546. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId })))
  547. expect(expectOk(await api.workspace.list(request({}))).archivedSessionIds).toEqual([])
  548. const abort = new AbortController()
  549. const stream: AsyncIterator<RpcRequest<HostFrame>> =
  550. api.events.host(request({}), abort.signal)[Symbol.asyncIterator]()
  551. const changed = nextHostFrame(stream)
  552. expect(expectOk(await api.workspace.archiveSession(request({ sessionId }))).archivedSessionIds)
  553. .toEqual([sessionId])
  554. expect(await changed).toMatchObject({
  555. payload: { type: 'host/archived-sessions-changed', archivedSessionIds: [sessionId] },
  556. })
  557. // Accounting and the session itself are untouched; list re-baselines the set.
  558. const listed = expectOk(await api.workspace.list(request({})))
  559. expect(listed.archivedSessionIds).toEqual([sessionId])
  560. expect(listed.items[0]?.sessionIds).toEqual([sessionId])
  561. expect(expectOk(await api.sessions.list(request({}))).items.map(item => item.sessionId)).toContain(sessionId)
  562. // The idempotent repeat emits no second frame: the next observed frame is
  563. // the workspace-changed of a later attach, not another archive snapshot.
  564. const after = nextHostFrame(stream)
  565. expect(expectOk(await api.workspace.archiveSession(request({ sessionId }))).archivedSessionIds)
  566. .toEqual([sessionId])
  567. const otherSession = SessionId('session-after-archive')
  568. expectOk(await api.sessions.create(request({ workspaceId: workspace.workspaceId, sessionId: otherSession })))
  569. expect((await after).payload.type).not.toBe('host/archived-sessions-changed')
  570. const missing = await api.workspace.archiveSession(request({ sessionId: SessionId('session-ghost') }))
  571. expect(missing.result).toMatchObject({
  572. ok: false,
  573. error: { code: 'session-not-found', details: { sessionId: 'session-ghost' } },
  574. })
  575. abort.abort()
  576. })
  577. })