workspace-controller.host.spec.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333
  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 { TypertRemoteFailure } 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. const roots: Context[] = []
  17. afterEach(async () => {
  18. await Promise.all(roots.splice(0).map(ctx => ctx.fiber.dispose()))
  19. })
  20. interface Deferred<T> {
  21. readonly promise: Promise<T>
  22. resolve(value: T): void
  23. }
  24. function deferred<T>(): Deferred<T> {
  25. let resolve!: (value: T) => void
  26. const promise = new Promise<T>((settle) => { resolve = settle })
  27. return { promise, resolve }
  28. }
  29. async function harness() {
  30. const root = realpathSync.native(mkdtempSync(join(tmpdir(), 'dsh-workspace-controller-')))
  31. const ctx = new Context()
  32. roots.push(ctx)
  33. await ctx.plugin(SessionStore)
  34. await ctx.plugin(Storage)
  35. ctx.storage.backend.register('memory', new MemoryStorageBackend())
  36. const storageDomain = new DomainFacility(ctx, { backend: 'memory', routes: {} })
  37. ctx.storage.mount('domain', storageDomain)
  38. ctx.provide('storageDomain', storageDomain)
  39. ctx.provide('sessionPersistence', { list: () => Promise.resolve([]) } as never)
  40. await ctx.plugin(WorkspaceRegistry)
  41. const dispose = (): void => {}
  42. ctx.provide('typert', {
  43. lookups: { configure: () => dispose },
  44. contexts: { configureHost: () => dispose },
  45. } as never)
  46. const controller = new WorkspaceController(ctx)
  47. return { controller, ctx, root, storageDomain }
  48. }
  49. function stageDir(root: string, name: string): string {
  50. const path = join(root, name)
  51. mkdirSync(path, { recursive: true })
  52. return path
  53. }
  54. async function nextFrame(
  55. iterator: AsyncIterator<WorkspaceFollowFrame>,
  56. ): Promise<WorkspaceFollowFrame> {
  57. const next = await iterator.next()
  58. if (next.done === true) throw new Error('Workspace stream ended before the expected frame')
  59. return next.value
  60. }
  61. describe('WorkspaceController commands', () => {
  62. it('serializes concurrent path adoption and preserves an existing title', async () => {
  63. const { controller, root } = await harness()
  64. const path = stageDir(root, 'alpha')
  65. const results = await Promise.all([
  66. controller.create({ path }),
  67. controller.create({ path }),
  68. ])
  69. const created = results.find(result => result.created)
  70. const resolved = results.find(result => !result.created)
  71. expect(created).toMatchObject({ workspace: { path, title: 'alpha' } })
  72. expect(resolved?.workspace.workspaceId).toBe(created?.workspace.workspaceId)
  73. const workspaceId = created?.workspace.workspaceId
  74. if (workspaceId === undefined) throw new Error('fixture did not create a Workspace')
  75. await controller.rename({ workspaceId, title: 'renamed' })
  76. await expect(controller.create({ path })).resolves.toMatchObject({
  77. created: false,
  78. workspace: { workspaceId, title: 'renamed' },
  79. })
  80. })
  81. it('maps invalid paths, blank names, conflicts, and unknown ids to stable failures', async () => {
  82. const { controller, root } = await harness()
  83. const first = await controller.create({ path: stageDir(root, 'first') })
  84. const second = await controller.create({ path: stageDir(root, 'second') })
  85. await expect(controller.create({ path: join(root, 'missing') })).rejects.toMatchObject({
  86. failure: { code: 'workspace-invalid-path', details: { path: join(root, 'missing') } },
  87. })
  88. expect(existsSync(join(root, 'missing'))).toBe(false)
  89. await expect(controller.rename({ workspaceId: first.workspace.workspaceId, title: ' ' }))
  90. .rejects.toMatchObject({ failure: { code: 'bad-request' } })
  91. await controller.rename({ workspaceId: first.workspace.workspaceId, title: 'occupied' })
  92. await expect(controller.rename({ workspaceId: second.workspace.workspaceId, title: ' occupied ' }))
  93. .rejects.toMatchObject({ failure: { code: 'workspace-name-conflict' } })
  94. await expect(controller.delete({ workspaceId: 'missing' as WorkspaceId }))
  95. .rejects.toMatchObject({ failure: { code: 'workspace-not-found' } })
  96. })
  97. it('preserves Remote failures and propagates unexpected registry failures', async () => {
  98. const { controller, ctx, root } = await harness()
  99. const remoteFailure = new TypertRemoteFailure({
  100. code: 'fixture-failure',
  101. message: 'already mapped',
  102. details: {},
  103. })
  104. const resolveByPath = vi.spyOn(ctx.workspaceRegistry, 'resolveByPath')
  105. .mockRejectedValueOnce(remoteFailure)
  106. .mockRejectedValueOnce('plain failure')
  107. await expect(controller.create({ path: stageDir(root, 'remote-failure') }))
  108. .rejects.toBe(remoteFailure)
  109. const plainFailure = controller.create({ path: stageDir(root, 'plain-failure') })
  110. await expect(plainFailure).rejects.toMatchObject({
  111. failure: { code: 'workspace-invalid-path' },
  112. })
  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({ failure: { 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({ failure: { 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({ failure: { 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. failure: {
  187. code: 'workspace-move-invalid',
  188. details: { beforeSessionId: 'missing-anchor' },
  189. },
  190. })
  191. await expect(controller.insertSessionBefore({
  192. workspaceId: 'missing' as WorkspaceId,
  193. sessionId: session.id,
  194. })).rejects.toMatchObject({ failure: { code: 'workspace-not-found' } })
  195. await expect(controller.archiveSession({ sessionId: session.id }))
  196. .resolves.toEqual({ archivedSessionIds: [session.id] })
  197. await expect(controller.archiveSession({ sessionId: SessionId('unknown') }))
  198. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  199. })
  200. })
  201. describe('WorkspaceController follow', () => {
  202. it('seeds a new feed from existing rows and rejects an inconsistent registry commit', async () => {
  203. const { ctx, root } = await harness()
  204. const existing = await ctx.workspaceRegistry.create(stageDir(root, 'existing'))
  205. const feed = new WorkspaceFeed(ctx)
  206. expect(feed.baseline()).toMatchObject({
  207. items: [{ workspaceId: existing.id }],
  208. })
  209. expect(() => {
  210. ctx.emit('domain/changed', {
  211. domain: 'workspace',
  212. table: '',
  213. key: '',
  214. operation: 'put',
  215. value: {
  216. initialized: true,
  217. workspaceIds: ['missing'],
  218. archivedSessionIds: [],
  219. },
  220. })
  221. }).toThrow('references missing Workspace "missing"')
  222. })
  223. it('starts with a complete baseline and emits committed increments in domain order', async () => {
  224. const { controller, ctx, root } = await harness()
  225. const abort = new AbortController()
  226. const iterator = controller.follow(abort.signal)[Symbol.asyncIterator]()
  227. await expect(nextFrame(iterator)).resolves.toEqual({
  228. type: 'baseline',
  229. value: { items: [], archivedSessionIds: [] },
  230. })
  231. const first = await controller.create({ path: stageDir(root, 'first') })
  232. await expect(nextFrame(iterator)).resolves.toMatchObject({
  233. type: 'upsert', workspace: { workspaceId: first.workspace.workspaceId },
  234. })
  235. await expect(nextFrame(iterator)).resolves.toEqual({
  236. type: 'order', workspaceIds: [first.workspace.workspaceId],
  237. })
  238. await controller.rename({ workspaceId: first.workspace.workspaceId, title: 'renamed' })
  239. await expect(nextFrame(iterator)).resolves.toMatchObject({
  240. type: 'upsert', workspace: { title: 'renamed' },
  241. })
  242. const second = await controller.create({ path: stageDir(root, 'second') })
  243. await expect(nextFrame(iterator)).resolves.toMatchObject({
  244. type: 'upsert', workspace: { workspaceId: second.workspace.workspaceId },
  245. })
  246. await expect(nextFrame(iterator)).resolves.toEqual({
  247. type: 'order', workspaceIds: [second.workspace.workspaceId, first.workspace.workspaceId],
  248. })
  249. await controller.insertBefore({
  250. workspaceId: first.workspace.workspaceId,
  251. beforeWorkspaceId: second.workspace.workspaceId,
  252. })
  253. await expect(nextFrame(iterator)).resolves.toEqual({
  254. type: 'order',
  255. workspaceIds: [first.workspace.workspaceId, second.workspace.workspaceId],
  256. })
  257. const session = ctx.sessions.create(SessionId('archived'), {
  258. meta: { cwd: first.workspace.path },
  259. })
  260. await controller.archiveSession({ sessionId: session.id })
  261. await expect(nextFrame(iterator)).resolves.toEqual({
  262. type: 'archived', archivedSessionIds: [session.id],
  263. })
  264. await controller.delete({ workspaceId: second.workspace.workspaceId })
  265. await expect(nextFrame(iterator)).resolves.toEqual({
  266. type: 'order', workspaceIds: [first.workspace.workspaceId],
  267. })
  268. await expect(nextFrame(iterator)).resolves.toEqual({
  269. type: 'remove', workspaceId: second.workspace.workspaceId,
  270. })
  271. abort.abort()
  272. await expect(iterator.next()).resolves.toEqual({ done: true, value: undefined })
  273. })
  274. it('ignores unrelated domain writes and closes active followers on disposal', async () => {
  275. const { controller, ctx, root } = await harness()
  276. const abort = new AbortController()
  277. const iterator = controller.follow(abort.signal)[Symbol.asyncIterator]()
  278. await nextFrame(iterator)
  279. ctx.emit('domain/changed', {
  280. domain: 'other', table: 'records', key: 'x', operation: 'put', value: {},
  281. })
  282. ctx.emit('domain/changed', {
  283. domain: 'workspace', table: '', key: '', operation: 'deleted',
  284. })
  285. ctx.emit('domain/changed', {
  286. domain: 'workspace', table: 'other', key: 'x', operation: 'put', value: {},
  287. })
  288. ctx.emit('domain/changed', {
  289. domain: 'workspace', table: 'workspaces', key: 'unknown', operation: 'deleted',
  290. })
  291. const pending = iterator.next()
  292. const created = await controller.create({ path: stageDir(root, 'visible') })
  293. await expect(pending).resolves.toMatchObject({ value: { type: 'upsert' } })
  294. await expect(iterator.next()).resolves.toEqual({
  295. done: false,
  296. value: { type: 'order', workspaceIds: [created.workspace.workspaceId] },
  297. })
  298. const closing = iterator.next()
  299. await ctx.fiber.dispose()
  300. roots.splice(roots.indexOf(ctx), 1)
  301. await expect(closing).resolves.toEqual({ done: true, value: undefined })
  302. })
  303. })