transport.client.spec.ts 17 KB

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