transport.client.spec.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373
  1. /**
  2. * Workspace Controller client plugin, state stream, and command facade driven
  3. * through the assembled Gateway client: every `workspace/*` call crosses the
  4. * roster's own Connection and is answered by endpoint name.
  5. */
  6. import { describe, expect, onTestFinished, vi } from 'vitest'
  7. import { RemoteStreamCarrierError, type ClientRemote } from '@deepseek-ai/dsh-api-gateway/client'
  8. import { SessionId } from '@deepseek-ai/dsh-session/types'
  9. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  10. import { frames, openStream, type RemoteMock, type StreamScript } from '@deepseek-ai/dsh-remote-mock'
  11. import { createClientTest, type TestClient, webApp } from '@deepseek-ai/dsh-client-test-runtime/src/assembly/index.ts'
  12. import {
  13. ClientWorkspaceModel,
  14. createWorkspaceStateStream,
  15. WorkspaceController,
  16. WorkspaceCreateError,
  17. type WorkspaceFollowSink,
  18. } from '../src/client/index.ts'
  19. import type { WorkspaceFollowFrame, WorkspaceId } from '../src/types.ts'
  20. import { FOLLOW, baseline, err, followGenerations, workspace, workspaceWorld } from './remote/workspace.client.ts'
  21. const SELF = '@deepseek-ai/dsh-api-workspace-controller'
  22. /** The plugin as the web bundle composes it: itself plus the Gateway client, the Connection, and the Typert registry. */
  23. const PLUGIN_ROSTER = webApp.closure([SELF])
  24. /** A stream or model built by hand talks through the Gateway client alone. */
  25. const API_ROSTER = webApp.closure(['@deepseek-ai/dsh-api-gateway'])
  26. const pluginTest = createClientTest({ roster: PLUGIN_ROSTER })
  27. const it = createClientTest({ roster: API_ROSTER })
  28. /** The first client boot pays the cold module transform of the plugin cone. */
  29. const COLD_BOOT_TIMEOUT_MS = 60_000
  30. const wid = (id: string): WorkspaceId => id as WorkspaceId
  31. const sid = (id: string): SessionId => SessionId(id)
  32. /** Boot the plugin with `follow` answering its state stream. */
  33. async function pluginClient(mock: RemoteMock, start: () => Promise<TestClient>, follow: StreamScript): Promise<TestClient> {
  34. mock.stream(FOLLOW, follow)
  35. return start()
  36. }
  37. /** Boot the Gateway client with the Workspace command answers registered, so `remote.workspace` is provided. */
  38. async function gatewayClient(mock: RemoteMock, start: () => Promise<TestClient>): Promise<{ remote: ClientRemote; client: TestClient }> {
  39. mock.load(workspaceWorld)
  40. const client = await start()
  41. return { remote: client.ctx.remote, client }
  42. }
  43. const carrierLoss = (message: string): StreamScript => (_args, stream) => {
  44. stream.fail(new RemoteStreamCarrierError(message))
  45. }
  46. function accepts(overrides: Partial<WorkspaceFollowSink> = {}): WorkspaceFollowSink {
  47. const ignore = (): void => {}
  48. return {
  49. replaceBaseline: ignore,
  50. upsertView: ignore,
  51. removeView: ignore,
  52. replaceOrder: ignore,
  53. replaceArchived: ignore,
  54. ...overrides,
  55. }
  56. }
  57. function streamStates(mock: RemoteMock): string[] {
  58. return mock.log.streams(FOLLOW).map(row => row.state)
  59. }
  60. describe('Workspace Controller Client apply', () => {
  61. pluginTest('provides the Workspace service and stops its follow generation with the plugin fiber', async ({ mock, start }) => {
  62. const client = await pluginClient(mock, start, openStream([baseline('mounted')]))
  63. await vi.waitFor(() => {
  64. expect(client.ctx.workspaces.list.getSnapshot()).toMatchObject({
  65. phase: 'ready',
  66. state: 'idle',
  67. items: [{ workspaceId: 'mounted' }],
  68. })
  69. })
  70. await client.unload(SELF)
  71. expect(streamStates(client.mock)).toEqual(['cancelled'])
  72. expect(client.ctx.get('workspaces')).toBeUndefined()
  73. }, COLD_BOOT_TIMEOUT_MS)
  74. pluginTest('reopens the follow and re-provides the service across a Loader rebuild', async ({ mock, start }) => {
  75. const client = await pluginClient(mock, start, openStream([baseline('mounted')]))
  76. await vi.waitFor(() => {
  77. expect(client.ctx.workspaces.list.getSnapshot()).toMatchObject({ phase: 'ready', items: [{ workspaceId: 'mounted' }] })
  78. })
  79. const before = client.ctx.workspaces
  80. await client.reload(SELF)
  81. expect(streamStates(client.mock)).toEqual(['cancelled', 'open'])
  82. expect(client.ctx.workspaces).not.toBe(before)
  83. await vi.waitFor(() => {
  84. expect(client.ctx.workspaces.list.getSnapshot()).toMatchObject({ phase: 'ready', items: [{ workspaceId: 'mounted' }] })
  85. })
  86. })
  87. pluginTest('publishes exhausted carrier retries as a gateway/internal error state', async ({ mock, start }) => {
  88. // Neither generation reaches an accepted baseline, so the retry budget runs
  89. // out and the escaping carrier failure crosses the stream boundary marked.
  90. const client = await pluginClient(mock, start, followGenerations([
  91. carrierLoss('generation lost'),
  92. carrierLoss('generation lost again'),
  93. ]))
  94. await vi.waitFor(() => {
  95. expect(client.ctx.workspaces.list.getSnapshot()).toMatchObject({
  96. state: 'error',
  97. error: { code: 'gateway/internal', message: 'generation lost again' },
  98. })
  99. })
  100. expect(streamStates(client.mock)).toEqual(['failed', 'failed'])
  101. })
  102. pluginTest('marks carrier loss while retrying and publishes a later protocol failure', async ({ mock, start }) => {
  103. const carrierFailure = vi.spyOn(ClientWorkspaceModel.prototype, 'handleCarrierFailure')
  104. const streamFailure = vi.spyOn(ClientWorkspaceModel.prototype, 'handleStreamFailure')
  105. onTestFinished(() => {
  106. carrierFailure.mockRestore()
  107. streamFailure.mockRestore()
  108. })
  109. const client = await pluginClient(mock, start, followGenerations([
  110. (_args, stream) => {
  111. stream.push(baseline('old'))
  112. stream.fail(new RemoteStreamCarrierError('generation lost'))
  113. },
  114. openStream([baseline('fresh'), baseline('duplicate')]),
  115. ]))
  116. await vi.waitFor(() => {
  117. expect(client.ctx.workspaces.list.getSnapshot()).toMatchObject({
  118. phase: 'ready',
  119. state: 'error',
  120. items: [{ workspaceId: 'fresh' }],
  121. error: { code: 'gateway/internal', message: 'Workspace state stream emitted more than one opening snapshot' },
  122. })
  123. })
  124. expect(carrierFailure).toHaveBeenCalledOnce()
  125. expect(streamFailure).toHaveBeenCalledOnce()
  126. })
  127. })
  128. describe('Workspace state stream', () => {
  129. it('delivers one baseline followed by increments', async ({ mock, start }) => {
  130. const { remote } = await gatewayClient(mock, start)
  131. const opening = baseline('one')
  132. const view = opening.value.items[0]!
  133. const increments: WorkspaceFollowFrame[] = [
  134. { type: 'upsert', workspace: view },
  135. { type: 'remove', workspaceId: view.workspaceId },
  136. { type: 'order', workspaceIds: [view.workspaceId] },
  137. { type: 'archived', archivedSessionIds: [sid('session-one')] },
  138. ]
  139. mock.stream(FOLLOW, openStream([opening, ...increments]))
  140. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  141. const upsertView = vi.fn<WorkspaceFollowSink['upsertView']>()
  142. const removeView = vi.fn<WorkspaceFollowSink['removeView']>()
  143. const replaceOrder = vi.fn<WorkspaceFollowSink['replaceOrder']>()
  144. const replaceArchived = vi.fn<WorkspaceFollowSink['replaceArchived']>()
  145. const stream = createWorkspaceStateStream(remote, {
  146. accept: accepts({ replaceBaseline, upsertView, removeView, replaceOrder, replaceArchived }),
  147. failed: vi.fn(),
  148. })
  149. stream.start()
  150. stream.start()
  151. await vi.waitFor(() => { expect(replaceArchived).toHaveBeenCalledOnce() })
  152. expect(replaceBaseline).toHaveBeenCalledWith(opening.value)
  153. expect(upsertView).toHaveBeenCalledWith(view)
  154. expect(removeView).toHaveBeenCalledWith(view.workspaceId)
  155. expect(replaceOrder).toHaveBeenCalledWith([view.workspaceId])
  156. expect(replaceArchived).toHaveBeenCalledWith(['session-one'])
  157. await stream.dispose()
  158. expect(streamStates(mock)).toEqual(['cancelled'])
  159. })
  160. it('retains the old state across carrier loss and applies the replacement baseline', async ({ mock, start }) => {
  161. const { remote } = await gatewayClient(mock, start)
  162. const carrier = new RemoteStreamCarrierError('socket lost')
  163. mock.stream(FOLLOW, followGenerations([
  164. (_args, stream) => {
  165. stream.push(baseline('old'))
  166. stream.fail(carrier)
  167. },
  168. openStream([baseline('fresh')]),
  169. ]))
  170. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  171. const carrierFailed = vi.fn()
  172. const failed = vi.fn()
  173. const stream = createWorkspaceStateStream(remote, {
  174. accept: accepts({ replaceBaseline }),
  175. carrierFailed,
  176. failed,
  177. })
  178. stream.start()
  179. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledTimes(2) })
  180. expect(replaceBaseline.mock.calls.map(([value]) => value.items[0]?.title)).toEqual(['old', 'fresh'])
  181. expect(carrierFailed).toHaveBeenCalledWith(carrier)
  182. expect(failed).not.toHaveBeenCalled()
  183. await stream.dispose()
  184. })
  185. it('classifies a normal end after the opening baseline as carrier loss', async ({ mock, start }) => {
  186. const { remote } = await gatewayClient(mock, start)
  187. mock.stream(FOLLOW, followGenerations([
  188. frames([baseline('old')]),
  189. openStream([baseline('fresh')]),
  190. ]))
  191. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  192. const carrierFailed = vi.fn()
  193. const stream = createWorkspaceStateStream(remote, {
  194. accept: accepts({ replaceBaseline }),
  195. carrierFailed,
  196. failed: vi.fn(),
  197. })
  198. stream.start()
  199. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledTimes(2) })
  200. expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
  201. message: 'Workspace state stream ended without a terminal result',
  202. })
  203. await stream.dispose()
  204. })
  205. it('suppresses callback failure after disposal begins', async ({ mock, start }) => {
  206. const { remote } = await gatewayClient(mock, start)
  207. mock.stream(FOLLOW, frames([baseline()]))
  208. const failed = vi.fn()
  209. let closing: Promise<void> | undefined
  210. const stream = createWorkspaceStateStream(remote, {
  211. accept: accepts({
  212. replaceBaseline: () => {
  213. closing = stream.dispose()
  214. throw new Error('disposed callback')
  215. },
  216. }),
  217. failed,
  218. })
  219. stream.start()
  220. await vi.waitFor(() => { expect(closing).toBeDefined() })
  221. await closing
  222. expect(failed).not.toHaveBeenCalled()
  223. })
  224. it.for([
  225. {
  226. name: 'an increment before the baseline',
  227. items: [{ type: 'remove', workspaceId: wid('one') }] as WorkspaceFollowFrame[],
  228. message: 'update before its opening snapshot',
  229. },
  230. {
  231. name: 'a duplicate baseline',
  232. items: [baseline(), baseline()] as WorkspaceFollowFrame[],
  233. message: 'more than one opening snapshot',
  234. },
  235. {
  236. name: 'a normal end before the baseline',
  237. items: [] as WorkspaceFollowFrame[],
  238. message: 'ended before its opening snapshot',
  239. },
  240. ])('reports $name as a terminal failure', async ({ items, message }, { mock, start }) => {
  241. const { remote } = await gatewayClient(mock, start)
  242. mock.stream(FOLLOW, frames(items))
  243. const failed = vi.fn()
  244. const stream = createWorkspaceStateStream(remote, { accept: accepts(), failed })
  245. stream.start()
  246. await vi.waitFor(() => { expect(failed).toHaveBeenCalledOnce() })
  247. const failure: unknown = failed.mock.calls[0]?.[0]
  248. expect(failure).toBeInstanceOf(Error)
  249. if (!(failure instanceof Error)) throw new Error('expected Workspace stream failure')
  250. expect(failure.message).toContain(message)
  251. // A protocol failure is terminal: no retry opens a second generation.
  252. expect(mock.log.requests(FOLLOW)).toHaveLength(1)
  253. await stream.dispose()
  254. })
  255. it('restarts a live generation without reporting cancellation as failure', async ({ mock, start }) => {
  256. const { remote } = await gatewayClient(mock, start)
  257. mock.stream(FOLLOW, followGenerations([
  258. openStream([baseline('first')]),
  259. openStream([baseline('second')]),
  260. ]))
  261. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  262. const failed = vi.fn()
  263. const stream = createWorkspaceStateStream(remote, {
  264. accept: accepts({ replaceBaseline }),
  265. failed,
  266. })
  267. stream.start()
  268. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledOnce() })
  269. stream.restart()
  270. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledTimes(2) })
  271. expect(failed).not.toHaveBeenCalled()
  272. expect(streamStates(mock)).toEqual(['cancelled', 'open'])
  273. await stream.dispose()
  274. })
  275. })
  276. describe('WorkspaceController', () => {
  277. it('publishes the model source and exposes successful Workspace commands', async ({ mock, start }) => {
  278. const { remote, client } = await gatewayClient(mock, start)
  279. const model = new ClientWorkspaceModel(remote.workspace)
  280. model.replaceBaseline({ items: [workspace('one')], archivedSessionIds: [] })
  281. const controller = new WorkspaceController(client.ctx, model)
  282. expect(controller.list).toBe(model)
  283. expect(client.ctx.workspaces.list).toBe(model)
  284. await expect(controller.create({ path: '/work/created' })).resolves.toMatchObject({ workspaceId: 'created' })
  285. await expect(controller.rename(wid('one'), 'renamed')).resolves.toMatchObject({ title: 'renamed' })
  286. await expect(controller.insertBefore(wid('one'))).resolves.toBeUndefined()
  287. await expect(controller.insertSessionBefore(wid('one'), sid('session'))).resolves.toMatchObject({
  288. sessionIds: ['session'],
  289. })
  290. await expect(controller.archiveSession(sid('session'))).resolves.toBeUndefined()
  291. await expect(controller.delete(wid('one'))).resolves.toBeUndefined()
  292. // Each command crosses the wire as one positional request object.
  293. expect(mock.log.requests('workspace/create')).toEqual([{ path: '/work/created' }])
  294. expect(mock.log.requests('workspace/rename')).toEqual([{ workspaceId: 'one', title: 'renamed' }])
  295. expect(mock.log.requests('workspace/insertBefore')).toEqual([{ workspaceId: 'one' }])
  296. expect(mock.log.requests('workspace/insertSessionBefore')).toEqual([{ workspaceId: 'one', sessionId: 'session' }])
  297. expect(mock.log.requests('workspace/archiveSession')).toEqual([{ sessionId: 'session' }])
  298. expect(mock.log.requests('workspace/delete')).toEqual([{ workspaceId: 'one' }])
  299. })
  300. it('maps generated business failures to the command facade errors', async ({ mock, start }) => {
  301. const { remote, client } = await gatewayClient(mock, start)
  302. const controller = new WorkspaceController(client.ctx, new ClientWorkspaceModel(remote.workspace))
  303. const missingWorkspace = new RemoteError('workspace/not-found', 'gone', { workspaceId: wid('missing') })
  304. const missingSession = new RemoteError('session/not-found', 'missing session', { sessionId: sid('session') })
  305. mock.remote.workspace.create.mockResolvedValueOnce(err(new RemoteError('workspace/invalid-path', 'missing path', { path: '/missing' })))
  306. const create = controller.create({ path: '/missing' })
  307. await expect(create).rejects.toBeInstanceOf(WorkspaceCreateError)
  308. await expect(create).rejects.toThrow('workspace create failed: workspace/invalid-path: missing path')
  309. mock.remote.workspace.rename.mockResolvedValueOnce(err(missingWorkspace))
  310. await expect(controller.rename(wid('missing'), 'name')).rejects.toThrow('workspace rename failed: workspace/not-found: gone')
  311. mock.remote.workspace.delete.mockResolvedValueOnce(err(missingWorkspace))
  312. await expect(controller.delete(wid('missing'))).rejects.toThrow('workspace delete failed: workspace/not-found: gone')
  313. mock.remote.workspace.insertBefore.mockResolvedValueOnce(err(missingWorkspace))
  314. await expect(controller.insertBefore(wid('missing'))).rejects.toThrow('workspace reorder failed: workspace/not-found: gone')
  315. mock.remote.workspace.archiveSession.mockResolvedValueOnce(err(missingSession))
  316. await expect(controller.archiveSession(sid('session')))
  317. .rejects.toThrow('workspace session archive failed: session/not-found: missing session')
  318. mock.remote.workspace.insertSessionBefore.mockResolvedValueOnce(err(new RemoteError(
  319. 'workspace/move-invalid', 'invalid move', { workspaceId: wid('missing'), sessionId: sid('session') },
  320. )))
  321. await expect(controller.insertSessionBefore(wid('missing'), sid('session')))
  322. .rejects.toThrow('workspace move failed: workspace/move-invalid: invalid move')
  323. })
  324. it('receives a carrier throw as the client\'s gateway/internal fold, never as a rejection', async ({ mock, start }) => {
  325. const { remote, client } = await gatewayClient(mock, start)
  326. const controller = new WorkspaceController(client.ctx, new ClientWorkspaceModel(remote.workspace))
  327. mock.remote.workspace.create.mockImplementation(() => Promise.reject(new Error('create wire down')))
  328. const create = controller.create({ path: '/work/created' })
  329. await expect(create).rejects.toBeInstanceOf(WorkspaceCreateError)
  330. await expect(create).rejects.toThrow(
  331. 'workspace create failed: gateway/internal: client api: workspace/create failed: create wire down',
  332. )
  333. expect(mock.log.calls('workspace/create').map(call => call.state)).toEqual(['failed'])
  334. })
  335. })