index.ts 4.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119
  1. /** Workspace-specific adapter for the Gateway-owned snapshot stream lifecycle. */
  2. import type { Context } from '@deepseek-ai/cordis'
  3. import {
  4. RemoteSnapshotStream,
  5. RemoteStreamCarrierError,
  6. type ClientRemote,
  7. } from '@deepseek-ai/dsh-api-gateway/client'
  8. import type { WorkspaceFollowFrame, WorkspaceFollowIncrement } from '../types.ts'
  9. import type { WorkspaceFollowSink } from './model.ts'
  10. import { ClientWorkspaceModel } from './model.ts'
  11. import { WorkspaceController } from './service.ts'
  12. export { ClientWorkspaceModel } from './model.ts'
  13. export type {
  14. WorkspaceFollowSink, WorkspaceListPhase, WorkspaceRemote, WorkspaceSnapshot,
  15. } from './model.ts'
  16. export { WorkspaceController, WorkspaceCreateError } from './service.ts'
  17. export type { IWorkspaces, WorkspaceSource } from './service.ts'
  18. export type { WorkspaceId, WorkspaceView } from '../types.ts'
  19. type WorkspaceBaselineFrame = Extract<WorkspaceFollowFrame, { type: 'baseline' }>
  20. /** Gateway-owned snapshot stream configured for Workspace state. */
  21. export type WorkspaceStateStream = RemoteSnapshotStream<
  22. WorkspaceBaselineFrame,
  23. WorkspaceFollowIncrement
  24. >
  25. declare module '@deepseek-ai/cordis' {
  26. interface Context {
  27. /** React-free Client Workspace state and commands. */
  28. workspaces: import('./service.ts').IWorkspaces
  29. }
  30. }
  31. /** Required Client Remote services. */
  32. export const inject = ['remote', 'remote.workspace']
  33. /**
  34. * Install Client Workspace state, commands, and reconnecting follow control.
  35. * @param ctx - Client root Context.
  36. */
  37. export function apply(ctx: Context): void {
  38. const model = new ClientWorkspaceModel(ctx.remote.workspace)
  39. new WorkspaceController(ctx, model)
  40. const control = createWorkspaceStateStream(ctx.remote, {
  41. accept: model,
  42. carrierFailed: () => { model.handleCarrierFailure() },
  43. failed: (error) => { model.handleStreamFailure(error) },
  44. })
  45. control.start()
  46. ctx.effect(
  47. () => async () => { await control.dispose() },
  48. 'workspace-controller.client.control',
  49. )
  50. }
  51. /** Domain sinks used by the Workspace state stream. */
  52. export interface WorkspaceStateStreamOptions {
  53. /** Destinations for decoded Workspace state operations. */
  54. readonly accept: WorkspaceFollowSink
  55. /** Observe a retryable carrier loss before reconnection. */
  56. readonly carrierFailed?: (error: RemoteStreamCarrierError) => void
  57. /** Publish a terminal business or protocol failure. */
  58. readonly failed: (error: unknown) => void
  59. }
  60. /**
  61. * Create the reconnecting Workspace state stream.
  62. * @param remote - Client Remote face carrying the Workspace namespace and the stream factory.
  63. * @param options - Workspace state destinations.
  64. * @returns an unstarted stream owned by the Client Workspace runtime.
  65. */
  66. export function createWorkspaceStateStream(
  67. remote: ClientRemote,
  68. options: WorkspaceStateStreamOptions,
  69. ): WorkspaceStateStream {
  70. const stream = remote.$stream<WorkspaceFollowFrame>({
  71. name: 'Workspace state stream',
  72. open: signal => remote.workspace.follow(signal),
  73. ended: accepted => accepted
  74. ? new RemoteStreamCarrierError('Workspace state stream ended without a terminal result')
  75. : new Error('Workspace state stream ended before its opening snapshot'),
  76. ...(options.carrierFailed === undefined ? {} : { carrierFailed: options.carrierFailed }),
  77. })
  78. return new RemoteSnapshotStream<WorkspaceBaselineFrame, WorkspaceFollowIncrement>(stream, {
  79. name: 'Workspace state stream',
  80. isSnapshot: (frame): frame is WorkspaceBaselineFrame => frame.type === 'baseline',
  81. replace: (frame) => { options.accept.replaceBaseline(frame.value) },
  82. update: (frame) => { acceptIncrement(options.accept, frame) },
  83. failed: options.failed,
  84. })
  85. }
  86. function acceptIncrement(accept: WorkspaceFollowSink, frame: WorkspaceFollowIncrement): void {
  87. switch (frame.type) {
  88. case 'upsert':
  89. accept.upsertView(frame.workspace)
  90. return
  91. case 'remove':
  92. accept.removeView(frame.workspaceId)
  93. return
  94. case 'order':
  95. accept.replaceOrder(frame.workspaceIds)
  96. return
  97. case 'archived':
  98. accept.replaceArchived(frame.archivedSessionIds)
  99. return
  100. /* v8 ignore next -- the generated Remote codec validates this closed union */
  101. default:
  102. return assertNever(frame)
  103. }
  104. }
  105. /* v8 ignore next 3 -- closed-union backstop after generated Remote validation */
  106. function assertNever(value: never): never {
  107. throw new Error(`unreachable Workspace increment: ${JSON.stringify(value)}`)
  108. }