workspace-controller.host.spec.ts 14 KB

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