transport.client.spec.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494
  1. import { Context } from '@deepseek-ai/cordis'
  2. import { describe, expect, it, vi } from 'vitest'
  3. import {
  4. RemoteStream,
  5. RemoteStreamCarrierError,
  6. type ClientRemote,
  7. type RemoteStreamOptions,
  8. } from '@deepseek-ai/dsh-api-gateway/client'
  9. import type { ConnectionHandle } from '@deepseek-ai/dsh-client-connection/client'
  10. import { SessionId } from '@deepseek-ai/dsh-session/types'
  11. import { RemoteError, type RemoteFailure, type RemoteResult } from '@deepseek-ai/dsh-typert-protocol'
  12. import * as WorkspaceClientPlugin from '../src/client/index.ts'
  13. import {
  14. ClientWorkspaceModel,
  15. createWorkspaceStateStream,
  16. WorkspaceController,
  17. WorkspaceCreateError,
  18. type WorkspaceFollowSink,
  19. type WorkspaceRemote,
  20. } from '../src/client/index.ts'
  21. import type {
  22. WorkspaceArchiveSessionRequest,
  23. WorkspaceArchiveValue,
  24. WorkspaceCreateRequest,
  25. WorkspaceCreateValue,
  26. WorkspaceDeleteRequest,
  27. WorkspaceDeleteValue,
  28. WorkspaceFollowFrame,
  29. WorkspaceInsertBeforeRequest,
  30. WorkspaceInsertSessionBeforeRequest,
  31. WorkspaceOrderValue,
  32. WorkspaceRenameRequest,
  33. WorkspaceId,
  34. WorkspaceValue,
  35. WorkspaceView,
  36. } from '../src/types.ts'
  37. const AVAILABLE_CONNECTION = {
  38. generation: {
  39. getSnapshot: () => ({ id: 1, host: { home: '/home/fixture' } }),
  40. subscribe: () => () => {},
  41. },
  42. }
  43. function workspaceClient(
  44. remote: WorkspaceRemote,
  45. connection: Pick<ConnectionHandle, 'generation'> = AVAILABLE_CONNECTION,
  46. ): ClientRemote {
  47. return {
  48. workspace: remote,
  49. $stream: <Item>(options: RemoteStreamOptions<Item>) => new RemoteStream(connection, options),
  50. } as unknown as ClientRemote
  51. }
  52. interface Generation {
  53. readonly frames: readonly WorkspaceFollowFrame[]
  54. readonly error?: unknown
  55. readonly hold?: boolean
  56. readonly afterAbort?: () => void
  57. readonly afterAbortError?: unknown
  58. }
  59. const baseline = (id?: string): Extract<WorkspaceFollowFrame, { type: 'baseline' }> => ({
  60. type: 'baseline',
  61. value: {
  62. items: id === undefined ? [] : [{
  63. workspaceId: id as never,
  64. path: `/work/${id}`,
  65. title: id,
  66. sessionIds: [],
  67. createdAt: '2026-01-01T00:00:00.000Z',
  68. updatedAt: '2026-01-01T00:00:00.000Z',
  69. }],
  70. archivedSessionIds: [],
  71. },
  72. })
  73. const wid = (id: string): WorkspaceId => id as WorkspaceId
  74. const sid = (id: string): SessionId => SessionId(id)
  75. function workspace(id: string, overrides: Partial<WorkspaceView> = {}): WorkspaceView {
  76. return {
  77. workspaceId: wid(id),
  78. path: `/work/${id}`,
  79. title: id,
  80. sessionIds: [],
  81. createdAt: '2026-01-01T00:00:00.000Z',
  82. updatedAt: '2026-01-01T00:00:00.000Z',
  83. ...overrides,
  84. }
  85. }
  86. function remoteOk<T>(value: T): RemoteResult<T> {
  87. return { ok: true, value }
  88. }
  89. function remoteFailure(error: RemoteFailure): RemoteResult<never> {
  90. return { ok: false, error }
  91. }
  92. function accepts(overrides: Partial<WorkspaceFollowSink> = {}): WorkspaceFollowSink {
  93. const ignore = (): void => {}
  94. return {
  95. replaceBaseline: ignore,
  96. upsertView: ignore,
  97. removeView: ignore,
  98. replaceOrder: ignore,
  99. replaceArchived: ignore,
  100. ...overrides,
  101. }
  102. }
  103. class ScriptedWorkspaceRemote implements WorkspaceRemote {
  104. readonly signals: AbortSignal[] = []
  105. calls = 0
  106. constructor(private readonly generations: readonly Generation[]) {}
  107. create(_request: WorkspaceCreateRequest): Promise<RemoteResult<WorkspaceCreateValue>> {
  108. throw new Error('unused')
  109. }
  110. rename(_request: WorkspaceRenameRequest): Promise<RemoteResult<WorkspaceValue>> {
  111. throw new Error('unused')
  112. }
  113. delete(_request: WorkspaceDeleteRequest): Promise<RemoteResult<WorkspaceDeleteValue>> {
  114. throw new Error('unused')
  115. }
  116. insertBefore(_request: WorkspaceInsertBeforeRequest): Promise<RemoteResult<WorkspaceOrderValue>> {
  117. throw new Error('unused')
  118. }
  119. insertSessionBefore(_request: WorkspaceInsertSessionBeforeRequest): Promise<RemoteResult<WorkspaceValue>> {
  120. throw new Error('unused')
  121. }
  122. archiveSession(_request: WorkspaceArchiveSessionRequest): Promise<RemoteResult<WorkspaceArchiveValue>> {
  123. throw new Error('unused')
  124. }
  125. async *follow(signal = new AbortController().signal): AsyncIterable<WorkspaceFollowFrame> {
  126. const generation = this.generations[this.calls++]
  127. if (generation === undefined) throw new Error('no scripted Workspace generation')
  128. this.signals.push(signal)
  129. for (const frame of generation.frames) yield frame
  130. if (generation.error !== undefined) throw generation.error
  131. if (generation.hold === true && !signal.aborted) {
  132. await new Promise<void>((resolve) => {
  133. signal.addEventListener('abort', () => { resolve() }, { once: true })
  134. })
  135. generation.afterAbort?.()
  136. if (generation.afterAbortError !== undefined) throw generation.afterAbortError
  137. }
  138. }
  139. }
  140. class CommandWorkspaceRemote implements WorkspaceRemote {
  141. readonly create = vi.fn<WorkspaceRemote['create']>(request => Promise.resolve(remoteOk({
  142. workspace: workspace('created', { path: request.path }),
  143. created: true,
  144. })))
  145. readonly rename = vi.fn<WorkspaceRemote['rename']>(request => Promise.resolve(remoteOk({
  146. workspace: workspace(String(request.workspaceId), { title: request.title }),
  147. })))
  148. readonly delete = vi.fn<WorkspaceRemote['delete']>(() => Promise.resolve(remoteOk({ deleted: true })))
  149. readonly insertBefore = vi.fn<WorkspaceRemote['insertBefore']>(request => Promise.resolve(remoteOk({
  150. workspaceIds: [request.workspaceId],
  151. })))
  152. readonly insertSessionBefore = vi.fn<WorkspaceRemote['insertSessionBefore']>(request => Promise.resolve(remoteOk({
  153. workspace: workspace(String(request.workspaceId), { sessionIds: [request.sessionId] }),
  154. })))
  155. readonly archiveSession = vi.fn<WorkspaceRemote['archiveSession']>(request => Promise.resolve(remoteOk({
  156. archivedSessionIds: [request.sessionId],
  157. })))
  158. async *follow(_signal?: AbortSignal): AsyncIterable<WorkspaceFollowFrame> {}
  159. }
  160. async function waitFor(check: () => void): Promise<void> {
  161. for (let attempt = 0; attempt < 40; attempt++) {
  162. try {
  163. check()
  164. return
  165. } catch {
  166. await Promise.resolve()
  167. }
  168. }
  169. check()
  170. }
  171. function provideClientServices(ctx: Context, remote: WorkspaceRemote): void {
  172. const connection: ConnectionHandle = {
  173. isLoopback: true,
  174. generation: AVAILABLE_CONNECTION.generation,
  175. state: { getSnapshot: () => 'connected' as const, subscribe: () => () => {} },
  176. rpc: {
  177. call: () => Promise.reject(new Error('unexpected generic RPC call')),
  178. },
  179. reconnect: () => {},
  180. registerGenerationSource: () => () => {},
  181. start: () => ({ stop: () => {} }),
  182. }
  183. ctx.reflect.provide('connection', connection)
  184. ctx.reflect.provide('remote', workspaceClient(remote, connection))
  185. ctx.reflect.provide('remote.workspace', remote)
  186. }
  187. describe('Workspace Controller Client apply', () => {
  188. it('provides the Workspace service and stops its follow generation with the plugin fiber', async () => {
  189. const ctx = new Context()
  190. const remote = new ScriptedWorkspaceRemote([{ frames: [baseline('mounted')], hold: true }])
  191. provideClientServices(ctx, remote)
  192. const fiber = ctx.plugin(WorkspaceClientPlugin)
  193. await fiber
  194. await waitFor(() => {
  195. expect(ctx.workspaces.list.getSnapshot()).toMatchObject({
  196. phase: 'ready',
  197. state: 'idle',
  198. items: [{ workspaceId: 'mounted' }],
  199. })
  200. })
  201. await fiber.dispose()
  202. expect(remote.signals[0]?.aborted).toBe(true)
  203. expect(ctx.get('workspaces')).toBeUndefined()
  204. })
  205. it('publishes exhausted carrier retries as a gateway/internal error state', async () => {
  206. const ctx = new Context()
  207. // Neither generation reaches an accepted baseline, so the retry budget runs
  208. // out and the escaping carrier failure crosses the stream boundary marked.
  209. const remote = new ScriptedWorkspaceRemote([
  210. { frames: [], error: new RemoteStreamCarrierError('generation lost') },
  211. { frames: [], error: new RemoteStreamCarrierError('generation lost again') },
  212. ])
  213. provideClientServices(ctx, remote)
  214. const fiber = ctx.plugin(WorkspaceClientPlugin)
  215. await fiber
  216. await waitFor(() => {
  217. expect(ctx.workspaces.list.getSnapshot()).toMatchObject({
  218. state: 'error',
  219. error: { code: 'gateway/internal', message: 'generation lost again' },
  220. })
  221. })
  222. expect(remote.calls).toBe(2)
  223. await fiber.dispose()
  224. })
  225. it('marks carrier loss while retrying and publishes a later protocol failure', async () => {
  226. const ctx = new Context()
  227. const remote = new ScriptedWorkspaceRemote([
  228. {
  229. frames: [baseline('old')],
  230. error: new RemoteStreamCarrierError('generation lost'),
  231. },
  232. { frames: [baseline('fresh'), baseline('duplicate')] },
  233. ])
  234. provideClientServices(ctx, remote)
  235. const carrierFailure = vi.spyOn(ClientWorkspaceModel.prototype, 'handleCarrierFailure')
  236. const streamFailure = vi.spyOn(ClientWorkspaceModel.prototype, 'handleStreamFailure')
  237. const fiber = ctx.plugin(WorkspaceClientPlugin)
  238. await fiber
  239. await waitFor(() => {
  240. expect(ctx.workspaces.list.getSnapshot()).toMatchObject({
  241. phase: 'ready',
  242. state: 'error',
  243. items: [{ workspaceId: 'fresh' }],
  244. error: { code: 'gateway/internal', message: 'Workspace state stream emitted more than one opening snapshot' },
  245. })
  246. })
  247. expect(carrierFailure).toHaveBeenCalledOnce()
  248. expect(streamFailure).toHaveBeenCalledOnce()
  249. await fiber.dispose()
  250. })
  251. })
  252. describe('Workspace state stream', () => {
  253. it('delivers one baseline followed by increments', async () => {
  254. const opening = baseline('one')
  255. const workspace = opening.value.items[0]!
  256. const remote = new ScriptedWorkspaceRemote([{
  257. frames: [
  258. opening,
  259. { type: 'upsert', workspace },
  260. { type: 'remove', workspaceId: workspace.workspaceId },
  261. { type: 'order', workspaceIds: [workspace.workspaceId] },
  262. { type: 'archived', archivedSessionIds: ['session-one' as never] },
  263. ],
  264. hold: true,
  265. }])
  266. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  267. const upsertView = vi.fn<WorkspaceFollowSink['upsertView']>()
  268. const removeView = vi.fn<WorkspaceFollowSink['removeView']>()
  269. const replaceOrder = vi.fn<WorkspaceFollowSink['replaceOrder']>()
  270. const replaceArchived = vi.fn<WorkspaceFollowSink['replaceArchived']>()
  271. const accept = accepts({
  272. replaceBaseline,
  273. upsertView,
  274. removeView,
  275. replaceOrder,
  276. replaceArchived,
  277. })
  278. const stream = createWorkspaceStateStream(workspaceClient(remote), {
  279. accept,
  280. failed: vi.fn(),
  281. })
  282. stream.start()
  283. stream.start()
  284. await vi.waitFor(() => { expect(replaceArchived).toHaveBeenCalledOnce() })
  285. expect(replaceBaseline).toHaveBeenCalledWith(opening.value)
  286. expect(upsertView).toHaveBeenCalledWith(workspace)
  287. expect(removeView).toHaveBeenCalledWith(workspace.workspaceId)
  288. expect(replaceOrder).toHaveBeenCalledWith([workspace.workspaceId])
  289. expect(replaceArchived).toHaveBeenCalledWith(['session-one'])
  290. await stream.dispose()
  291. expect(remote.signals[0]?.aborted).toBe(true)
  292. })
  293. it('retains the old state across carrier loss and applies the replacement baseline', async () => {
  294. const carrier = new RemoteStreamCarrierError('socket lost')
  295. const remote = new ScriptedWorkspaceRemote([
  296. { frames: [baseline('old')], error: carrier },
  297. { frames: [baseline('fresh')], hold: true },
  298. ])
  299. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  300. const carrierFailed = vi.fn()
  301. const failed = vi.fn()
  302. const stream = createWorkspaceStateStream(workspaceClient(remote), {
  303. accept: accepts({ replaceBaseline }),
  304. carrierFailed,
  305. failed,
  306. })
  307. stream.start()
  308. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledTimes(2) })
  309. expect(replaceBaseline.mock.calls.map(([value]) => value.items[0]?.title)).toEqual(['old', 'fresh'])
  310. expect(carrierFailed).toHaveBeenCalledWith(carrier)
  311. expect(failed).not.toHaveBeenCalled()
  312. await stream.dispose()
  313. })
  314. it('classifies a normal end after the opening baseline as carrier loss', async () => {
  315. const remote = new ScriptedWorkspaceRemote([
  316. { frames: [baseline('old')] },
  317. { frames: [baseline('fresh')], hold: true },
  318. ])
  319. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  320. const carrierFailed = vi.fn()
  321. const stream = createWorkspaceStateStream(workspaceClient(remote), {
  322. accept: accepts({ replaceBaseline }),
  323. carrierFailed,
  324. failed: vi.fn(),
  325. })
  326. stream.start()
  327. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledTimes(2) })
  328. expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
  329. message: 'Workspace state stream ended without a terminal result',
  330. })
  331. await stream.dispose()
  332. })
  333. it('suppresses callback failure after disposal begins', async () => {
  334. const failed = vi.fn()
  335. let closing: Promise<void> | undefined
  336. const stream = createWorkspaceStateStream(
  337. workspaceClient(new ScriptedWorkspaceRemote([{ frames: [baseline()] }])),
  338. {
  339. accept: accepts({
  340. replaceBaseline: () => {
  341. closing = stream.dispose()
  342. throw new Error('disposed callback')
  343. },
  344. }),
  345. failed,
  346. },
  347. )
  348. stream.start()
  349. await vi.waitFor(() => { expect(closing).toBeDefined() })
  350. await closing
  351. expect(failed).not.toHaveBeenCalled()
  352. })
  353. it.each([
  354. {
  355. name: 'an increment before the baseline',
  356. frames: [{ type: 'remove', workspaceId: 'one' as never }] as WorkspaceFollowFrame[],
  357. message: 'update before its opening snapshot',
  358. },
  359. {
  360. name: 'a duplicate baseline',
  361. frames: [baseline(), baseline()] as WorkspaceFollowFrame[],
  362. message: 'more than one opening snapshot',
  363. },
  364. {
  365. name: 'a normal end before the baseline',
  366. frames: [] as WorkspaceFollowFrame[],
  367. message: 'ended before its opening snapshot',
  368. },
  369. ])('reports $name as a terminal failure', async ({ frames, message }) => {
  370. const failed = vi.fn()
  371. const stream = createWorkspaceStateStream(
  372. workspaceClient(new ScriptedWorkspaceRemote([{ frames }])),
  373. { accept: accepts(), failed },
  374. )
  375. stream.start()
  376. await vi.waitFor(() => { expect(failed).toHaveBeenCalledOnce() })
  377. const failure: unknown = failed.mock.calls[0]?.[0]
  378. expect(failure).toBeInstanceOf(Error)
  379. if (!(failure instanceof Error)) throw new Error('expected Workspace stream failure')
  380. expect(failure.message).toContain(message)
  381. await stream.dispose()
  382. })
  383. it('restarts a live generation without reporting cancellation as failure', async () => {
  384. const remote = new ScriptedWorkspaceRemote([
  385. { frames: [baseline('first')], hold: true },
  386. { frames: [baseline('second')], hold: true },
  387. ])
  388. const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
  389. const failed = vi.fn()
  390. const stream = createWorkspaceStateStream(workspaceClient(remote), {
  391. accept: accepts({ replaceBaseline }),
  392. failed,
  393. })
  394. stream.start()
  395. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledOnce() })
  396. stream.restart()
  397. await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledTimes(2) })
  398. expect(failed).not.toHaveBeenCalled()
  399. await stream.dispose()
  400. })
  401. })
  402. describe('WorkspaceController', () => {
  403. it('publishes the model source and exposes successful Workspace commands', async () => {
  404. const remote = new CommandWorkspaceRemote()
  405. const model = new ClientWorkspaceModel(remote)
  406. model.replaceBaseline({ items: [workspace('one')], archivedSessionIds: [] })
  407. const controller = new WorkspaceController(new Context(), model)
  408. expect(controller.list).toBe(model)
  409. await expect(controller.create({ path: '/work/created' })).resolves.toMatchObject({ workspaceId: 'created' })
  410. await expect(controller.rename(wid('one'), 'renamed')).resolves.toMatchObject({ title: 'renamed' })
  411. await expect(controller.insertBefore(wid('one'))).resolves.toBeUndefined()
  412. await expect(controller.insertSessionBefore(wid('one'), sid('session'))).resolves.toMatchObject({
  413. sessionIds: ['session'],
  414. })
  415. await expect(controller.archiveSession(sid('session'))).resolves.toBeUndefined()
  416. await expect(controller.delete(wid('one'))).resolves.toBeUndefined()
  417. })
  418. it('maps generated business failures to the command facade errors', async () => {
  419. const remote = new CommandWorkspaceRemote()
  420. const controller = new WorkspaceController(new Context(), new ClientWorkspaceModel(remote))
  421. const missingWorkspace = new RemoteError('workspace/not-found', 'gone', { workspaceId: wid('missing') })
  422. const missingSession = new RemoteError('session/not-found', 'missing session', { sessionId: sid('session') })
  423. remote.create.mockResolvedValueOnce(remoteFailure(new RemoteError('workspace/invalid-path', 'missing path', { path: '/missing' })))
  424. const create = controller.create({ path: '/missing' })
  425. await expect(create).rejects.toBeInstanceOf(WorkspaceCreateError)
  426. await expect(create).rejects.toThrow('workspace/invalid-path: missing path')
  427. remote.rename.mockResolvedValueOnce(remoteFailure(missingWorkspace))
  428. await expect(controller.rename(wid('missing'), 'name')).rejects.toThrow('workspace rename failed: workspace/not-found: gone')
  429. remote.delete.mockResolvedValueOnce(remoteFailure(missingWorkspace))
  430. await expect(controller.delete(wid('missing'))).rejects.toThrow('workspace delete failed: workspace/not-found: gone')
  431. remote.insertBefore.mockResolvedValueOnce(remoteFailure(missingWorkspace))
  432. await expect(controller.insertBefore(wid('missing'))).rejects.toThrow('workspace reorder failed: workspace/not-found: gone')
  433. remote.archiveSession.mockResolvedValueOnce(remoteFailure(missingSession))
  434. await expect(controller.archiveSession(sid('session')))
  435. .rejects.toThrow('workspace session archive failed: session/not-found: missing session')
  436. remote.insertSessionBefore.mockResolvedValueOnce(remoteFailure(new RemoteError(
  437. 'workspace/move-invalid', 'invalid move', { workspaceId: wid('missing'), sessionId: sid('session') },
  438. )))
  439. await expect(controller.insertSessionBefore(wid('missing'), sid('session')))
  440. .rejects.toThrow('workspace move failed: workspace/move-invalid: invalid move')
  441. })
  442. })