fake-api.ts 6.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163
  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, HostFrame, IApiClient, MuxFrame, RpcError, RpcReceipt, RpcRequest, RpcResponse, SessionId,
  6. } from '@deepseek-ai/dsh-client-connection/client'
  7. import { RpcId } from '@deepseek-ai/dsh-client-connection/client'
  8. export interface Deferred<T> {
  9. promise: Promise<T>
  10. resolve(value: T): void
  11. reject(error: unknown): void
  12. }
  13. /** Test-held settlement: the case decides when an RPC lands (history-pending injections etc.). */
  14. export function deferred<T>(): Deferred<T> {
  15. let resolve!: (value: T) => void
  16. let reject!: (error: unknown) => void
  17. const promise = new Promise<T>((res, rej) => {
  18. resolve = res
  19. reject = rej
  20. })
  21. return { promise, resolve, reject }
  22. }
  23. let nextRpc = 0
  24. export function ok<T>(value: T): RpcResponse<T> {
  25. return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: true, value } }
  26. }
  27. export function err<T>(error: RpcError): RpcResponse<T> {
  28. return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: false, error } }
  29. }
  30. type StreamItem<F> = { kind: 'frame'; envelope: RpcRequest<F> } | { kind: 'end' } | { kind: 'fail'; error: unknown }
  31. interface StreamConn<F> {
  32. feed(item: StreamItem<F>): void
  33. }
  34. export class FakeApiClient implements IApiClient {
  35. /** Chronological call record: [method, payload]. */
  36. readonly calls: { method: string; payload: unknown }[] = []
  37. // Programmable slots (defaults answer OK-empty); reassign per case.
  38. onList: (payload: unknown) => Promise<RpcResponse<{ items: never[] }>> = () => Promise.resolve(ok({ items: [] }))
  39. onCreate: (payload: unknown) => Promise<RpcResponse<{ sessionId: SessionId }>> = () => Promise.resolve(ok({ sessionId: 'fk-new' as SessionId }))
  40. onHistory: (payload: { sessionId: SessionId; beforeSeq?: number; maxMessages?: number })
  41. => Promise<RpcResponse<{ events: never[]; hasMore: boolean }>> =
  42. () => Promise.resolve(ok({ events: [], hasMore: false }))
  43. onPrompt: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  44. onCancel: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  45. onDescribe: (payload: unknown) => Promise<RpcResponse<{ version: string; cwd: string; attachedSessions: number }>> =
  46. () => Promise.resolve(ok({ version: '0-fake', cwd: '/f', attachedSessions: 0 }))
  47. private readonly muxConns: StreamConn<MuxFrame>[] = []
  48. private readonly hostConns: StreamConn<HostFrame>[] = []
  49. // Parameters carry local structural annotations: the CI lint lane runs
  50. // without built lib/, so IApiClient's indexed-access types collapse to any
  51. // and inferred parameters would trip no-unsafe-argument.
  52. readonly sessions: IApiClient['sessions'] = {
  53. list: (payload: unknown) => this.record('session.list', payload, this.onList(payload)),
  54. create: (payload: unknown) => this.record('session.create', payload, this.onCreate(payload)),
  55. history: (payload: { sessionId: SessionId; beforeSeq?: number; maxMessages?: number }) =>
  56. this.record('session.history', payload, this.onHistory(payload)),
  57. prompt: (payload: unknown) => this.record('session.prompt', payload, this.onPrompt(payload)),
  58. cancel: (payload: unknown) => this.record('session.cancel', payload, this.onCancel(payload)),
  59. }
  60. readonly host: IApiClient['host'] = {
  61. describe: (payload: unknown) => this.record('host.describe', payload, this.onDescribe(payload)),
  62. }
  63. /** When true, streams never fire onOpen (misbehaving-carrier material for the handshake timeout guard). */
  64. suppressStreamOpen = false
  65. /** When true, onOpen callbacks are parked instead of fired; releaseStreamOpens() fires them.
  66. * Lets a case hold the readiness handshake open (describe done, streams not yet "established"). */
  67. holdStreamOpen = false
  68. private heldOpens: (() => void)[] = []
  69. releaseStreamOpens(): void {
  70. const held = this.heldOpens
  71. this.heldOpens = []
  72. for (const fire of held) fire()
  73. }
  74. readonly events: IApiClient['events'] = {
  75. mux: (_payload: unknown, signal: AbortSignal, onOpen?: () => void) => this.openStream(this.muxConns, signal, onOpen),
  76. host: (_payload: unknown, signal: AbortSignal, onOpen?: () => void) => this.openStream(this.hostConns, signal, onOpen),
  77. }
  78. onRespond: (message: ClientResponse) => Promise<RpcReceipt> = () => Promise.resolve({ accepted: true })
  79. respond(message: ClientResponse): Promise<RpcReceipt> {
  80. return this.record('respond', message, this.onRespond(message))
  81. }
  82. /** Push one mux frame to every open mux stream (rpcId minted unless pinned by the case). */
  83. pushMux(frame: MuxFrame, rpcId?: string): void {
  84. for (const conn of [...this.muxConns]) conn.feed({ kind: 'frame', envelope: { rpcId: RpcId(rpcId ?? `push-${nextRpc++}`), payload: frame } })
  85. }
  86. pushHost(frame: HostFrame, rpcId?: string): void {
  87. for (const conn of [...this.hostConns]) conn.feed({ kind: 'frame', envelope: { rpcId: RpcId(rpcId ?? `push-${nextRpc++}`), payload: frame } })
  88. }
  89. /** End (clean close) or fail (throw) every open stream — reconnect-path material. */
  90. endStreams(): void {
  91. for (const conn of [...this.muxConns, ...this.hostConns]) conn.feed({ kind: 'end' })
  92. }
  93. failStreams(error: unknown): void {
  94. for (const conn of [...this.muxConns, ...this.hostConns]) conn.feed({ kind: 'fail', error })
  95. }
  96. get openMuxCount(): number {
  97. return this.muxConns.length
  98. }
  99. callsOf(method: string): unknown[] {
  100. return this.calls.filter(c => c.method === method).map(c => c.payload)
  101. }
  102. private record<T>(method: string, payload: unknown, response: Promise<T>): Promise<T> {
  103. this.calls.push({ method, payload })
  104. return response
  105. }
  106. private async *openStream<F>(registry: StreamConn<F>[], signal: AbortSignal, onOpen?: () => void): AsyncGenerator<RpcRequest<F>> {
  107. const inbox: StreamItem<F>[] = []
  108. let wake: (() => void) | null = null
  109. const conn: StreamConn<F> = {
  110. feed: (item) => {
  111. inbox.push(item)
  112. wake?.()
  113. },
  114. }
  115. registry.push(conn)
  116. if (this.holdStreamOpen && onOpen !== undefined) this.heldOpens.push(onOpen)
  117. else if (!this.suppressStreamOpen) onOpen?.()
  118. try {
  119. while (!signal.aborted) {
  120. while (inbox.length > 0) {
  121. const item = inbox.shift() as StreamItem<F>
  122. if (item.kind === 'end') return
  123. if (item.kind === 'fail') throw item.error
  124. yield item.envelope
  125. }
  126. await new Promise<void>((resolve) => {
  127. wake = resolve
  128. signal.addEventListener('abort', () => { resolve() }, { once: true })
  129. })
  130. wake = null
  131. }
  132. } finally {
  133. registry.splice(registry.indexOf(conn), 1)
  134. }
  135. }
  136. }