fake-api.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236
  1. // Test-local programmable IApiClient fake (NOT the fixture: fixture is a demo
  2. // data source on a real clock; behavior tests need per-case responses and
  3. // deferred-controlled timing). Streams are hand pumps: pushMux/pushHost.
  4. import type {
  5. ClientResponse, CommandDescriptor, CommandExecuteResult, HostFrame, IApiClient, ModelTarget, MuxFrame,
  6. RpcError, RpcReceipt, RpcRequest, RpcResponse, SessionId, SessionModels, SkillEntry,
  7. WorkspaceId, WorkspaceView,
  8. } from '@deepseek-ai/dsh-client-connection/client'
  9. import { RpcId } from '@deepseek-ai/dsh-client-connection/client'
  10. /** Programmable-default workspace row (branded id, ISO-ish times). */
  11. function fakeWorkspace(id: string, over: Partial<WorkspaceView> = {}): WorkspaceView {
  12. return {
  13. workspaceId: id as WorkspaceId,
  14. path: '/f/ws',
  15. title: 'ws',
  16. sessionIds: [],
  17. createdAt: '2026-01-01T00:00:00.000Z',
  18. updatedAt: '2026-01-01T00:00:00.000Z',
  19. ...over,
  20. }
  21. }
  22. export interface Deferred<T> {
  23. promise: Promise<T>
  24. resolve(value: T): void
  25. reject(error: unknown): void
  26. }
  27. /** Test-held settlement: the case decides when an RPC lands (history-pending injections etc.). */
  28. export function deferred<T>(): Deferred<T> {
  29. let resolve!: (value: T) => void
  30. let reject!: (error: unknown) => void
  31. const promise = new Promise<T>((res, rej) => {
  32. resolve = res
  33. reject = rej
  34. })
  35. return { promise, resolve, reject }
  36. }
  37. let nextRpc = 0
  38. export function ok<T>(value: T): RpcResponse<T> {
  39. return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: true, value } }
  40. }
  41. export function err<T>(error: RpcError): RpcResponse<T> {
  42. return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: false, error } }
  43. }
  44. type StreamItem<F> = { kind: 'frame'; envelope: RpcRequest<F> } | { kind: 'end' } | { kind: 'fail'; error: unknown }
  45. interface StreamConn<F> {
  46. feed(item: StreamItem<F>): void
  47. }
  48. export class FakeApiClient implements IApiClient {
  49. /** Chronological call record: [method, payload]. */
  50. readonly calls: { method: string; payload: unknown }[] = []
  51. // Programmable slots (defaults answer OK-empty); reassign per case.
  52. onList: (payload: unknown) => Promise<RpcResponse<{ items: never[] }>> = () => Promise.resolve(ok({ items: [] }))
  53. onCreate: (payload: unknown) => Promise<RpcResponse<{ sessionId: SessionId }>> = () => Promise.resolve(ok({ sessionId: 'fk-new' as SessionId }))
  54. readonly defaultModel: ModelTarget = { provider: 'deepseek', model: 'deepseek-v4-flash' }
  55. onHistory: (payload: { sessionId: SessionId; beforeSeq?: number; maxMessages?: number })
  56. => Promise<RpcResponse<{ events: never[]; hasMore: boolean; todos?: { content: string; status: 'pending' | 'in_progress' | 'completed' }[] }>> =
  57. () => Promise.resolve(ok({ events: [], hasMore: false }))
  58. onModels: (payload: unknown) => Promise<RpcResponse<SessionModels>> = () => Promise.resolve(ok({
  59. current: this.defaultModel,
  60. groups: [{
  61. id: 'deepseek',
  62. name: 'DeepSeek',
  63. models: [{ id: 'deepseek-v4-flash', name: 'DeepSeek V4 Flash' }],
  64. }],
  65. failures: [],
  66. }))
  67. onSelectModel: (payload: { provider: string; model: string }) =>
  68. Promise<RpcResponse<{ selected: ModelTarget }>> =
  69. payload => Promise.resolve(ok({ selected: { provider: payload.provider, model: payload.model } }))
  70. onPrompt: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  71. onCancel: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  72. onDescribe: (payload: unknown) => Promise<RpcResponse<{ version: string; cwd: string; attachedSessions: number }>> =
  73. () => Promise.resolve(ok({ version: '0-fake', cwd: '/f', attachedSessions: 0 }))
  74. onPickDirectory: (payload: unknown) => Promise<RpcResponse<{ path: string | null }>> =
  75. () => Promise.resolve(ok({ path: null }))
  76. private readonly muxConns: StreamConn<MuxFrame>[] = []
  77. private readonly hostConns: StreamConn<HostFrame>[] = []
  78. // Parameters carry local structural annotations: the CI lint lane runs
  79. // without built lib/, so IApiClient's indexed-access types collapse to any
  80. // and inferred parameters would trip no-unsafe-argument.
  81. readonly sessions: IApiClient['sessions'] = {
  82. list: (payload: unknown) => this.record('session.list', payload, this.onList(payload)),
  83. create: (payload: unknown) => this.record('session.create', payload, this.onCreate(payload)),
  84. history: (payload: { sessionId: SessionId; beforeSeq?: number; maxMessages?: number }) =>
  85. this.record('session.history', payload, this.onHistory(payload)),
  86. models: (payload: unknown) => this.record('session.models', payload, this.onModels(payload)),
  87. selectModel: (payload: { provider: string; model: string }) =>
  88. this.record('session.selectModel', payload, this.onSelectModel(payload)),
  89. prompt: (payload: unknown) => this.record('session.prompt', payload, this.onPrompt(payload)),
  90. cancel: (payload: unknown) => this.record('session.cancel', payload, this.onCancel(payload)),
  91. }
  92. readonly host: IApiClient['host'] = {
  93. describe: (payload: unknown) => this.record('host.describe', payload, this.onDescribe(payload)),
  94. pickDirectory: (payload: unknown) => this.record('host.pickDirectory', payload, this.onPickDirectory(payload)),
  95. }
  96. onWorkspaceList: (payload: unknown) => Promise<RpcResponse<{ items: never[] }>> = () => Promise.resolve(ok({ items: [] }))
  97. onWorkspaceCreate: (payload: unknown) => Promise<RpcResponse<{ workspace: WorkspaceView; created: boolean }>> =
  98. () => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws'), created: true }))
  99. onWorkspaceRename: (payload: unknown) => Promise<RpcResponse<{ workspace: WorkspaceView }>> =
  100. () => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws') }))
  101. onWorkspaceDelete: (payload: unknown) => Promise<RpcResponse<{ deleted: true }>> =
  102. () => Promise.resolve(ok({ deleted: true }))
  103. onWorkspaceInsertSessionBefore: (payload: unknown) => Promise<RpcResponse<{ workspace: WorkspaceView }>> =
  104. () => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws') }))
  105. readonly workspace: IApiClient['workspace'] = {
  106. list: (payload: unknown) => this.record('workspace.list', payload, this.onWorkspaceList(payload)),
  107. create: (payload: unknown) => this.record('workspace.create', payload, this.onWorkspaceCreate(payload)),
  108. rename: (payload: unknown) => this.record('workspace.rename', payload, this.onWorkspaceRename(payload)),
  109. delete: (payload: unknown) => this.record('workspace.delete', payload, this.onWorkspaceDelete(payload)),
  110. insertSessionBefore: (payload: unknown) =>
  111. this.record('workspace.insertSessionBefore', payload, this.onWorkspaceInsertSessionBefore(payload)),
  112. }
  113. // Payloads stay `unknown` (lint-lane note above); response rows are the real
  114. // wire shapes so cases can program requires-bearing catalogs and dual-address
  115. // skill lists without casts.
  116. onCommandList: (payload: unknown) => Promise<RpcResponse<{ commands: CommandDescriptor[] }>> = () => Promise.resolve(ok({ commands: [] }))
  117. onCommandExecute: (payload: unknown) => Promise<RpcResponse<{ matched: boolean; result?: CommandExecuteResult }>> =
  118. () => Promise.resolve(ok({ matched: false }))
  119. onSkillList: (payload: unknown) => Promise<RpcResponse<{ skills: SkillEntry[] }>> = () => Promise.resolve(ok({ skills: [] }))
  120. readonly commands: IApiClient['commands'] = {
  121. list: (payload: unknown) => this.record('command.list', payload, this.onCommandList(payload)),
  122. execute: (payload: unknown) => this.record('command.execute', payload, this.onCommandExecute(payload)),
  123. }
  124. readonly skills: IApiClient['skills'] = {
  125. list: (payload: unknown) => this.record('skill.list', payload, this.onSkillList(payload)),
  126. }
  127. /** When true, streams never fire onOpen (misbehaving-carrier material for the handshake timeout guard). */
  128. suppressStreamOpen = false
  129. /** When true, onOpen callbacks are parked instead of fired; releaseStreamOpens() fires them.
  130. * Lets a case hold the readiness handshake open (describe done, streams not yet "established"). */
  131. holdStreamOpen = false
  132. private heldOpens: (() => void)[] = []
  133. releaseStreamOpens(): void {
  134. const held = this.heldOpens
  135. this.heldOpens = []
  136. for (const fire of held) fire()
  137. }
  138. readonly events: IApiClient['events'] = {
  139. mux: (_payload: unknown, signal: AbortSignal, onOpen?: () => void) => this.openStream(this.muxConns, signal, onOpen),
  140. host: (_payload: unknown, signal: AbortSignal, onOpen?: () => void) => this.openStream(this.hostConns, signal, onOpen),
  141. }
  142. onRespond: (message: ClientResponse) => Promise<RpcReceipt> = () => Promise.resolve({ accepted: true })
  143. respond(message: ClientResponse): Promise<RpcReceipt> {
  144. return this.record('respond', message, this.onRespond(message))
  145. }
  146. /** Push one mux frame to every open mux stream (rpcId minted unless pinned by the case). */
  147. pushMux(frame: MuxFrame, rpcId?: string): void {
  148. for (const conn of [...this.muxConns]) conn.feed({ kind: 'frame', envelope: { rpcId: RpcId(rpcId ?? `push-${nextRpc++}`), payload: frame } })
  149. }
  150. pushHost(frame: HostFrame, rpcId?: string): void {
  151. for (const conn of [...this.hostConns]) conn.feed({ kind: 'frame', envelope: { rpcId: RpcId(rpcId ?? `push-${nextRpc++}`), payload: frame } })
  152. }
  153. /** End (clean close) or fail (throw) every open stream — reconnect-path material. */
  154. endStreams(): void {
  155. for (const conn of [...this.muxConns, ...this.hostConns]) conn.feed({ kind: 'end' })
  156. }
  157. failStreams(error: unknown): void {
  158. for (const conn of [...this.muxConns, ...this.hostConns]) conn.feed({ kind: 'fail', error })
  159. }
  160. get openMuxCount(): number {
  161. return this.muxConns.length
  162. }
  163. callsOf(method: string): unknown[] {
  164. return this.calls.filter(c => c.method === method).map(c => c.payload)
  165. }
  166. private record<T>(method: string, payload: unknown, response: Promise<T>): Promise<T> {
  167. this.calls.push({ method, payload })
  168. return response
  169. }
  170. private async *openStream<F>(registry: StreamConn<F>[], signal: AbortSignal, onOpen?: () => void): AsyncGenerator<RpcRequest<F>> {
  171. const inbox: StreamItem<F>[] = []
  172. let wake: (() => void) | null = null
  173. const conn: StreamConn<F> = {
  174. feed: (item) => {
  175. inbox.push(item)
  176. wake?.()
  177. },
  178. }
  179. registry.push(conn)
  180. if (this.holdStreamOpen && onOpen !== undefined) this.heldOpens.push(onOpen)
  181. else if (!this.suppressStreamOpen) onOpen?.()
  182. try {
  183. while (!signal.aborted) {
  184. while (inbox.length > 0) {
  185. const item = inbox.shift() as StreamItem<F>
  186. if (item.kind === 'end') return
  187. if (item.kind === 'fail') throw item.error
  188. yield item.envelope
  189. }
  190. await new Promise<void>((resolve) => {
  191. wake = resolve
  192. signal.addEventListener('abort', () => { resolve() }, { once: true })
  193. })
  194. wake = null
  195. }
  196. } finally {
  197. registry.splice(registry.indexOf(conn), 1)
  198. }
  199. }
  200. }