fake-api.ts 9.6 KB

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