workspace-controller.host.spec.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332
  1. import { existsSync, mkdirSync, mkdtempSync, realpathSync } from 'node:fs'
  2. import { tmpdir } from 'node:os'
  3. import { join } from 'node:path'
  4. import { afterEach, describe, expect, it, vi } from 'vitest'
  5. import { Context } from '@deepseek-ai/cordis'
  6. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  7. import Storage from '@deepseek-ai/dsh-storage'
  8. import { DomainFacility } from '@deepseek-ai/dsh-storage-domain'
  9. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  10. import WorkspaceRegistry from '@deepseek-ai/dsh-workspace'
  11. import type { WorkspaceId } from '@deepseek-ai/dsh-workspace/types'
  12. import WorkspaceController from '../src/index.ts'
  13. import { WorkspaceFeed } from '../src/feed.ts'
  14. import type { WorkspaceFollowFrame } from '../src/types.ts'
  15. import { MemoryStorageBackend } from '../../../storage/storage-domain/tests/helpers/memory-backend.ts'
  16. declare module '@deepseek-ai/dsh-typert-protocol' {
  17. interface RemoteErrorDetailsMap {
  18. 'fixture/failure': {}
  19. }
  20. }
  21. const roots: Context[] = []
  22. afterEach(async () => {
  23. await Promise.all(roots.splice(0).map(ctx => ctx.fiber.dispose()))
  24. })
  25. interface Deferred<T> {
  26. readonly promise: Promise<T>
  27. resolve(value: T): void
  28. }
  29. function deferred<T>(): Deferred<T> {
  30. let resolve!: (value: T) => void
  31. const promise = new Promise<T>((settle) => { resolve = settle })
  32. return { promise, resolve }
  33. }
  34. async function harness() {
  35. const root = realpathSync.native(mkdtempSync(join(tmpdir(), 'dsh-workspace-controller-')))
  36. const ctx = new Context()
  37. roots.push(ctx)
  38. await ctx.plugin(SessionStore)
  39. await ctx.plugin(Storage)
  40. ctx.storage.backend.register('memory', new MemoryStorageBackend())
  41. const storageDomain = new DomainFacility(ctx, { backend: 'memory', routes: {} })
  42. ctx.storage.mount('domain', storageDomain)
  43. ctx.provide('storageDomain', storageDomain)
  44. ctx.provide('sessionPersistence', { list: () => Promise.resolve([]) } as never)
  45. await ctx.plugin(WorkspaceRegistry)
  46. const dispose = (): void => {}
  47. ctx.provide('typert', {
  48. lookups: { configure: () => dispose },
  49. contexts: { configureHost: () => dispose },
  50. } as never)
  51. const controller = new WorkspaceController(ctx)
  52. return { controller, ctx, root, storageDomain }
  53. }
  54. function stageDir(root: string, name: string): string {
  55. const path = join(root, name)
  56. mkdirSync(path, { recursive: true })
  57. return path
  58. }
  59. async function nextFrame(
  60. iterator: AsyncIterator<WorkspaceFollowFrame>,
  61. ): Promise<WorkspaceFollowFrame> {
  62. const next = await iterator.next()
  63. if (next.done === true) throw new Error('Workspace stream ended before the expected frame')
  64. return next.value
  65. }
  66. describe('WorkspaceController commands', () => {
  67. it('serializes concurrent path adoption and preserves an existing title', async () => {
  68. const { controller, root } = await harness()
  69. const path = stageDir(root, 'alpha')
  70. const results = await Promise.all([
  71. controller.create({ path }),
  72. controller.create({ path }),
  73. ])
  74. const created = results.find(result => result.created)
  75. const resolved = results.find(result => !result.created)
  76. expect(created).toMatchObject({ workspace: { path, title: 'alpha' } })
  77. expect(resolved?.workspace.workspaceId).toBe(created?.workspace.workspaceId)
  78. const workspaceId = created?.workspace.workspaceId
  79. if (workspaceId === undefined) throw new Error('fixture did not create a Workspace')
  80. await controller.rename({ workspaceId, title: 'renamed' })
  81. await expect(controller.create({ path })).resolves.toMatchObject({
  82. created: false,
  83. workspace: { workspaceId, title: 'renamed' },
  84. })
  85. })
  86. it('maps invalid paths, blank names, conflicts, and unknown ids to stable failures', async () => {
  87. const { controller, root } = await harness()
  88. const first = await controller.create({ path: stageDir(root, 'first') })
  89. const second = await controller.create({ path: stageDir(root, 'second') })
  90. await expect(controller.create({ path: join(root, 'missing') })).rejects.toMatchObject({
  91. code: 'workspace/invalid-path',
  92. details: { path: join(root, 'missing') },
  93. })
  94. expect(existsSync(join(root, 'missing'))).toBe(false)
  95. await expect(controller.rename({ workspaceId: first.workspace.workspaceId, title: ' ' }))
  96. .rejects.toMatchObject({ code: 'gateway/bad-request' })
  97. await controller.rename({ workspaceId: first.workspace.workspaceId, title: 'occupied' })
  98. await expect(controller.rename({ workspaceId: second.workspace.workspaceId, title: ' occupied ' }))
  99. .rejects.toMatchObject({ code: 'workspace/name-conflict' })
  100. await expect(controller.delete({ workspaceId: 'missing' as WorkspaceId }))
  101. .rejects.toMatchObject({ code: 'workspace/not-found' })
  102. })
  103. it('preserves Remote failures and propagates unexpected registry failures', async () => {
  104. const { controller, ctx, root } = await harness()
  105. const remoteFailure = new RemoteError('fixture/failure', 'already mapped', {})
  106. const resolveByPath = vi.spyOn(ctx.workspaceRegistry, 'resolveByPath')
  107. .mockRejectedValueOnce(remoteFailure)
  108. .mockRejectedValueOnce('plain failure')
  109. await expect(controller.create({ path: stageDir(root, 'remote-failure') }))
  110. .rejects.toBe(remoteFailure)
  111. const plainFailure = controller.create({ path: stageDir(root, 'plain-failure') })
  112. await expect(plainFailure).rejects.toMatchObject({ code: 'workspace/invalid-path' })
  113. await expect(plainFailure).rejects.toThrow('plain failure')
  114. resolveByPath.mockRestore()
  115. const created = await controller.create({ path: stageDir(root, 'created') })
  116. const workspace = ctx.workspaceRegistry.get(created.workspace.workspaceId)
  117. if (workspace === undefined) throw new Error('fixture Workspace disappeared')
  118. const orderFailure = new Error('order storage failed')
  119. vi.spyOn(ctx.workspaceRegistry, 'insertBefore').mockRejectedValueOnce(orderFailure)
  120. await expect(controller.insertBefore({ workspaceId: created.workspace.workspaceId }))
  121. .rejects.toBe(orderFailure)
  122. const moveFailure = new Error('membership storage failed')
  123. vi.spyOn(workspace, 'insertSessionBefore').mockRejectedValueOnce(moveFailure)
  124. await expect(controller.insertSessionBefore({
  125. workspaceId: created.workspace.workspaceId,
  126. sessionId: SessionId('session'),
  127. })).rejects.toBe(moveFailure)
  128. const archiveFailure = new Error('archive storage failed')
  129. vi.spyOn(ctx.workspaceRegistry, 'archiveSession').mockRejectedValueOnce(archiveFailure)
  130. await expect(controller.archiveSession({ sessionId: SessionId('session') }))
  131. .rejects.toBe(archiveFailure)
  132. })
  133. it('resolves queued Workspace identities when their operation starts', async () => {
  134. const { controller, ctx, root } = await harness()
  135. const target = await controller.create({ path: stageDir(root, 'target') })
  136. const blockerPath = stageDir(root, 'blocker')
  137. const gate = deferred<undefined>()
  138. const originalResolveByPath = ctx.workspaceRegistry.resolveByPath.bind(ctx.workspaceRegistry)
  139. const resolveByPath = vi.spyOn(ctx.workspaceRegistry, 'resolveByPath')
  140. resolveByPath.mockImplementationOnce(async (path) => {
  141. await gate.promise
  142. return originalResolveByPath(path)
  143. })
  144. const blocker = controller.create({ path: blockerPath })
  145. const deletion = controller.delete({ workspaceId: target.workspace.workspaceId })
  146. const staleRename = controller.rename({
  147. workspaceId: target.workspace.workspaceId,
  148. title: 'must-not-land',
  149. })
  150. gate.resolve(undefined)
  151. await blocker
  152. await expect(deletion).resolves.toEqual({ deleted: true })
  153. await expect(staleRename).rejects.toMatchObject({ code: 'workspace/not-found' })
  154. })
  155. it('reorders Workspaces and Sessions and archives only known Sessions', async () => {
  156. const { controller, ctx, root } = await harness()
  157. const first = await controller.create({ path: stageDir(root, 'first') })
  158. const second = await controller.create({ path: stageDir(root, 'second') })
  159. await expect(controller.insertBefore({
  160. workspaceId: first.workspace.workspaceId,
  161. beforeWorkspaceId: second.workspace.workspaceId,
  162. })).resolves.toEqual({
  163. workspaceIds: [first.workspace.workspaceId, second.workspace.workspaceId],
  164. })
  165. await expect(controller.insertBefore({ workspaceId: 'missing' as WorkspaceId }))
  166. .rejects.toMatchObject({ code: 'workspace/not-found' })
  167. const session = ctx.sessions.create(SessionId('session-one'), {
  168. meta: { cwd: first.workspace.path },
  169. })
  170. const workspace = ctx.workspaceRegistry.get(first.workspace.workspaceId)
  171. if (workspace === undefined) throw new Error('fixture Workspace disappeared')
  172. await workspace.attachSession(session.id)
  173. await expect(controller.insertSessionBefore({
  174. workspaceId: first.workspace.workspaceId,
  175. sessionId: session.id,
  176. })).resolves.toMatchObject({ workspace: { sessionIds: [session.id] } })
  177. await expect(controller.insertSessionBefore({
  178. workspaceId: first.workspace.workspaceId,
  179. sessionId: SessionId('missing-session'),
  180. })).rejects.toMatchObject({ code: 'workspace/move-invalid' })
  181. await expect(controller.insertSessionBefore({
  182. workspaceId: first.workspace.workspaceId,
  183. sessionId: session.id,
  184. beforeSessionId: SessionId('missing-anchor'),
  185. })).rejects.toMatchObject({
  186. code: 'workspace/move-invalid',
  187. details: { beforeSessionId: 'missing-anchor' },
  188. })
  189. await expect(controller.insertSessionBefore({
  190. workspaceId: 'missing' as WorkspaceId,
  191. sessionId: session.id,
  192. })).rejects.toMatchObject({ code: 'workspace/not-found' })
  193. await expect(controller.archiveSession({ sessionId: session.id }))
  194. .resolves.toEqual({ archivedSessionIds: [session.id] })
  195. await expect(controller.archiveSession({ sessionId: SessionId('unknown') }))
  196. .rejects.toMatchObject({ code: 'session/not-found' })
  197. })
  198. })
  199. describe('WorkspaceController follow', () => {
  200. it('seeds a new feed from existing rows and rejects an inconsistent registry commit', async () => {
  201. const { ctx, root } = await harness()
  202. const existing = await ctx.workspaceRegistry.create(stageDir(root, 'existing'))
  203. const feed = new WorkspaceFeed(ctx)
  204. expect(feed.baseline()).toMatchObject({
  205. items: [{ workspaceId: existing.id }],
  206. })
  207. expect(() => {
  208. ctx.emit('domain/changed', {
  209. domain: 'workspace',
  210. table: '',
  211. key: '',
  212. operation: 'put',
  213. value: {
  214. initialized: true,
  215. workspaceIds: ['missing'],
  216. archivedSessionIds: [],
  217. },
  218. })
  219. }).toThrow('references missing Workspace "missing"')
  220. })
  221. it('starts with a complete baseline and emits committed increments in domain order', async () => {
  222. const { controller, ctx, root } = await harness()
  223. const abort = new AbortController()
  224. const iterator = controller.follow(abort.signal)[Symbol.asyncIterator]()
  225. await expect(nextFrame(iterator)).resolves.toEqual({
  226. type: 'baseline',
  227. value: { items: [], archivedSessionIds: [] },
  228. })
  229. const first = await controller.create({ path: stageDir(root, 'first') })
  230. await expect(nextFrame(iterator)).resolves.toMatchObject({
  231. type: 'upsert', workspace: { workspaceId: first.workspace.workspaceId },
  232. })
  233. await expect(nextFrame(iterator)).resolves.toEqual({
  234. type: 'order', workspaceIds: [first.workspace.workspaceId],
  235. })
  236. await controller.rename({ workspaceId: first.workspace.workspaceId, title: 'renamed' })
  237. await expect(nextFrame(iterator)).resolves.toMatchObject({
  238. type: 'upsert', workspace: { title: 'renamed' },
  239. })
  240. const second = await controller.create({ path: stageDir(root, 'second') })
  241. await expect(nextFrame(iterator)).resolves.toMatchObject({
  242. type: 'upsert', workspace: { workspaceId: second.workspace.workspaceId },
  243. })
  244. await expect(nextFrame(iterator)).resolves.toEqual({
  245. type: 'order', workspaceIds: [second.workspace.workspaceId, first.workspace.workspaceId],
  246. })
  247. await controller.insertBefore({
  248. workspaceId: first.workspace.workspaceId,
  249. beforeWorkspaceId: second.workspace.workspaceId,
  250. })
  251. await expect(nextFrame(iterator)).resolves.toEqual({
  252. type: 'order',
  253. workspaceIds: [first.workspace.workspaceId, second.workspace.workspaceId],
  254. })
  255. const session = ctx.sessions.create(SessionId('archived'), {
  256. meta: { cwd: first.workspace.path },
  257. })
  258. await controller.archiveSession({ sessionId: session.id })
  259. await expect(nextFrame(iterator)).resolves.toEqual({
  260. type: 'archived', archivedSessionIds: [session.id],
  261. })
  262. await controller.delete({ workspaceId: second.workspace.workspaceId })
  263. await expect(nextFrame(iterator)).resolves.toEqual({
  264. type: 'order', workspaceIds: [first.workspace.workspaceId],
  265. })
  266. await expect(nextFrame(iterator)).resolves.toEqual({
  267. type: 'remove', workspaceId: second.workspace.workspaceId,
  268. })
  269. abort.abort()
  270. await expect(iterator.next()).resolves.toEqual({ done: true, value: undefined })
  271. })
  272. it('ignores unrelated domain writes and closes active followers on disposal', async () => {
  273. const { controller, ctx, root } = await harness()
  274. const abort = new AbortController()
  275. const iterator = controller.follow(abort.signal)[Symbol.asyncIterator]()
  276. await nextFrame(iterator)
  277. ctx.emit('domain/changed', {
  278. domain: 'other', table: 'records', key: 'x', operation: 'put', value: {},
  279. })
  280. ctx.emit('domain/changed', {
  281. domain: 'workspace', table: '', key: '', operation: 'deleted',
  282. })
  283. ctx.emit('domain/changed', {
  284. domain: 'workspace', table: 'other', key: 'x', operation: 'put', value: {},
  285. })
  286. ctx.emit('domain/changed', {
  287. domain: 'workspace', table: 'workspaces', key: 'unknown', operation: 'deleted',
  288. })
  289. const pending = iterator.next()
  290. const created = await controller.create({ path: stageDir(root, 'visible') })
  291. await expect(pending).resolves.toMatchObject({ value: { type: 'upsert' } })
  292. await expect(iterator.next()).resolves.toEqual({
  293. done: false,
  294. value: { type: 'order', workspaceIds: [created.workspace.workspaceId] },
  295. })
  296. const closing = iterator.next()
  297. await ctx.fiber.dispose()
  298. roots.splice(roots.indexOf(ctx), 1)
  299. await expect(closing).resolves.toEqual({ done: true, value: undefined })
  300. })
  301. })