fake-api.client.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497
  1. // Test-local programmable Remote 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). Session streams are hand pumps: pushFollow/pushControl.
  4. import type {
  5. MessageId,
  6. SessionId, SessionSearchItem,
  7. SubagentCatalog, SubagentInterruptReceipt, SubagentPromptReceipt,
  8. WorkspaceId, WorkspaceView,
  9. } from '@deepseek-ai/dsh-api-remotes/client'
  10. import type {
  11. SessionAddress,
  12. SessionAssistantStreamBaseline,
  13. SessionControlBaseline,
  14. SessionControlFrame,
  15. SessionEditRequest,
  16. SessionEditValue,
  17. SessionFollowFrame,
  18. SessionFollowRequest,
  19. SessionPage,
  20. SessionPageRequest,
  21. SessionProjectionBaseline,
  22. SessionSelectModelRequest,
  23. SessionSelectModelValue,
  24. } from '@deepseek-ai/dsh-api-session-controller/types'
  25. import type { WorkspaceRemote } from '@deepseek-ai/dsh-api-workspace-controller/client'
  26. import type { WorkspaceFollowFrame } from '@deepseek-ai/dsh-api-workspace-controller/types'
  27. import type { RemoteFailure, RemoteResult } from '@deepseek-ai/dsh-typert-protocol'
  28. import {
  29. RemoteStream,
  30. type RemoteStreamOptions,
  31. } from '@deepseek-ai/dsh-api-gateway/client'
  32. import type { SessionRemotes } from '../src/client/sessions/remotes.ts'
  33. import { historyRecordLastSeq } from '../src/client/sessions/history-records.ts'
  34. const AVAILABLE_STREAM_CONNECTION = {
  35. generation: {
  36. getSnapshot: () => ({ id: 1, host: { home: '/h' } }),
  37. subscribe: () => () => {},
  38. },
  39. }
  40. /** Programmable-default workspace row (branded id, ISO-ish times). */
  41. function fakeWorkspace(id: string, over: Partial<WorkspaceView> = {}): WorkspaceView {
  42. return {
  43. workspaceId: id as WorkspaceId,
  44. path: '/f/ws',
  45. title: 'ws',
  46. sessionIds: [],
  47. createdAt: '2026-01-01T00:00:00.000Z',
  48. updatedAt: '2026-01-01T00:00:00.000Z',
  49. ...over,
  50. }
  51. }
  52. function addressSessionId(address: SessionAddress): SessionId {
  53. return address.kind === 'session' ? address.sessionId : address.childSessionId
  54. }
  55. export interface Deferred<T> {
  56. promise: Promise<T>
  57. resolve(value: T): void
  58. reject(error: unknown): void
  59. }
  60. /** Test-held settlement: the case decides when an RPC lands (history-pending injections etc.). */
  61. export function deferred<T>(): Deferred<T> {
  62. let resolve!: (value: T) => void
  63. let reject!: (error: unknown) => void
  64. const promise = new Promise<T>((res, rej) => {
  65. resolve = res
  66. reject = rej
  67. })
  68. return { promise, resolve, reject }
  69. }
  70. /**
  71. * Successful generated Remote result for programmable domain fakes.
  72. * @param value - the value the Host answers with.
  73. * @returns the success branch of a Remote result.
  74. */
  75. export function ok<T>(value: T): RemoteResult<T> {
  76. return { ok: true, value }
  77. }
  78. /**
  79. * Failed generated Remote result carrying the owner's declared failure.
  80. * @param error - the owner-declared failure.
  81. * @returns the failure branch of a Remote result.
  82. */
  83. export function err<T>(error: RemoteFailure): RemoteResult<T> {
  84. return { ok: false, error }
  85. }
  86. type ValueStreamItem<F> =
  87. | { kind: 'frame'; value: F; delivered?: () => void }
  88. | { kind: 'end' }
  89. | { kind: 'fail'; error: unknown }
  90. interface ValueStreamConn<F> {
  91. feed(item: ValueStreamItem<F>): void
  92. }
  93. interface OpenValueStream<F> {
  94. readonly values: AsyncGenerator<F>
  95. dispose(): void
  96. }
  97. /**
  98. * Commands Remote double: the generated face delivers the carrier's outcome, so
  99. * a test that programs nothing sees an empty catalog and an unmatched line.
  100. * @returns the Remote namespaces the session cluster calls.
  101. */
  102. export type RuntimeRemotes = SessionRemotes & { readonly workspace: WorkspaceRemote }
  103. export function fakeRemote(api = new FakeApiClient()): RuntimeRemotes {
  104. return api.sessionRemotes()
  105. }
  106. export class FakeApiClient {
  107. /** Chronological call record: [method, payload]. */
  108. readonly calls: { method: string; payload: unknown }[] = []
  109. /** Session ids in physical follow-generation opening order. */
  110. readonly followStarts: SessionId[] = []
  111. // Programmable slots (defaults answer OK-empty); reassign per case.
  112. onList: (payload: unknown) => Promise<RemoteResult<{ items: never[] }>> = () => Promise.resolve(ok({ items: [] }))
  113. onSearch: (payload: unknown) => Promise<RemoteResult<{ items: SessionSearchItem[]; hasMore: boolean }>> =
  114. () => Promise.resolve(ok({ items: [], hasMore: false }))
  115. onCreate: (payload: unknown) => Promise<RemoteResult<{ sessionId: SessionId }>> = () => Promise.resolve(ok({ sessionId: 'fk-new' as SessionId }))
  116. onSelectModel: (payload: SessionSelectModelRequest) => Promise<RemoteResult<SessionSelectModelValue>> =
  117. payload => Promise.resolve(ok({
  118. selected: {
  119. provider: payload.provider,
  120. model: payload.model,
  121. ...(payload.reasoningEffort === undefined
  122. ? {}
  123. : { reasoningEffort: payload.reasoningEffort }),
  124. },
  125. }))
  126. onRename: (payload: unknown) => Promise<RemoteResult<{ title: string; seq: number }>> = () => Promise.resolve(ok({ title: 'fk-renamed', seq: 0 }))
  127. onFork: (payload: unknown) => Promise<RemoteResult<{ sessionId: SessionId }>> = () => Promise.resolve(ok({ sessionId: 'fk-fork' as SessionId }))
  128. onHistory: (payload: { sessionId: SessionId; throughSeq?: number; beforeSeq?: number; maxMessages?: number })
  129. => Promise<RemoteResult<SessionPage & { readonly projections?: SessionProjectionBaseline }>> =
  130. () => Promise.resolve(ok({ records: [], hasMore: false }))
  131. onPrompt: (payload: unknown) => Promise<RemoteResult<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  132. onEdit: (payload: SessionEditRequest) => Promise<RemoteResult<SessionEditValue>> = () =>
  133. Promise.resolve(ok({ accepted: true as const, messageSeq: 0 }))
  134. onAttachment: (payload: unknown) => Promise<RemoteResult<{ attachment: { attachmentId: never; mediaType: 'image/png'; bytes: number; width: number; height: number }; data: string }>> =
  135. () => Promise.resolve(ok({ attachment: { attachmentId: 'a' as never, mediaType: 'image/png', bytes: 1, width: 1, height: 1 }, data: 'AA==' }))
  136. onUpdateQueue: (payload: unknown) => Promise<RemoteResult<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  137. onCancel: (payload: unknown) => Promise<RemoteResult<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  138. onOpenWorkspacePath: (payload: unknown) => Promise<RemoteResult<{ opened: true }>> =
  139. () => Promise.resolve(ok({ opened: true as const }))
  140. private readonly followConns = new Map<SessionId, ValueStreamConn<SessionFollowFrame>[]>()
  141. private readonly controlConns: ValueStreamConn<SessionControlFrame>[] = []
  142. private readonly workspaceConns: ValueStreamConn<WorkspaceFollowFrame>[] = []
  143. /** Optional Host opening cursor override for stale-page and reconnect tests. */
  144. followCursor: number | undefined
  145. controlBaseline: SessionControlBaseline = {
  146. queues: {},
  147. jobs: {},
  148. projections: {},
  149. }
  150. assistantStreamBaseline: SessionAssistantStreamBaseline = {
  151. revision: 0,
  152. attempts: [],
  153. }
  154. workspaceBaseline: Extract<WorkspaceFollowFrame, { type: 'baseline' }>['value'] = {
  155. items: [],
  156. archivedSessionIds: [],
  157. }
  158. lastSearchSignal: AbortSignal | undefined
  159. onSubagentList: (payload: unknown) => Promise<RemoteResult<SubagentCatalog>>
  160. = () => Promise.resolve(ok({ entries: [], parentAvailable: true }))
  161. onSubagentPrompt: (payload: unknown) => Promise<RemoteResult<SubagentPromptReceipt>>
  162. = () => Promise.resolve(ok({ messageId: 'fake-message' as MessageId }))
  163. onSubagentInterrupt: (payload: unknown) => Promise<RemoteResult<SubagentInterruptReceipt>>
  164. = () => Promise.resolve(ok({ accepted: true as const }))
  165. onWorkspaceCreate: (payload: unknown) => Promise<RemoteResult<{ workspace: WorkspaceView; created: boolean }>> =
  166. () => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws'), created: true }))
  167. onWorkspaceRename: (payload: unknown) => Promise<RemoteResult<{ workspace: WorkspaceView }>> =
  168. () => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws') }))
  169. onWorkspaceDelete: (payload: unknown) => Promise<RemoteResult<{ deleted: true }>> =
  170. () => Promise.resolve(ok({ deleted: true }))
  171. onWorkspaceInsertBefore: (payload: unknown) => Promise<RemoteResult<{ workspaceIds: WorkspaceId[] }>> =
  172. () => Promise.resolve(ok({ workspaceIds: [] }))
  173. onWorkspaceInsertSessionBefore: (payload: unknown) => Promise<RemoteResult<{ workspace: WorkspaceView }>> =
  174. () => Promise.resolve(ok({ workspace: fakeWorkspace('fk-ws') }))
  175. onWorkspaceArchiveSession: (payload: unknown) => Promise<RemoteResult<{ archivedSessionIds: SessionId[] }>> =
  176. payload => Promise.resolve(ok({ archivedSessionIds: [(payload as { sessionId: SessionId }).sessionId] }))
  177. /** Remote namespaces bound to this fake's programmable unary slots and stream pumps. */
  178. sessionRemotes(): RuntimeRemotes {
  179. return {
  180. $stream: <Item>(options: RemoteStreamOptions<Item>) => (
  181. new RemoteStream(AVAILABLE_STREAM_CONNECTION, options)
  182. ),
  183. commands: {
  184. execute: () => Promise.resolve({ ok: true, value: undefined }),
  185. },
  186. session: {
  187. canOpenWorkspacePath: () => Promise.resolve(ok(true)),
  188. list: payload => this.record('session.list', payload, this.onList(payload)),
  189. modelCatalog: () => Promise.resolve({
  190. ok: true,
  191. value: {
  192. default: { provider: 'fixture', model: 'fixture' },
  193. routableProviders: [],
  194. groups: [],
  195. failures: [],
  196. },
  197. }),
  198. search: (payload, signal) => {
  199. this.lastSearchSignal = signal
  200. return this.record('session.search', payload, this.onSearch(payload))
  201. },
  202. create: payload => this.record('session.create', payload, this.onCreate(payload)),
  203. selectModel: payload => this.record(
  204. 'session.selectModel',
  205. payload,
  206. this.onSelectModel(payload),
  207. ),
  208. rename: payload => this.record('session.rename', payload, this.onRename(payload)),
  209. fork: payload => this.record('session.fork', payload, this.onFork(payload)),
  210. prompt: payload => this.record('session.prompt', payload, this.onPrompt(payload)),
  211. edit: payload => this.record('session.edit', payload, this.onEdit(payload)),
  212. attachment: payload => this.record('session.attachment', payload, this.onAttachment(payload)),
  213. updateQueue: payload => this.record('session.updateQueue', payload, this.onUpdateQueue(payload)),
  214. cancel: payload => this.record('session.cancel', payload, this.onCancel(payload)),
  215. openWorkspacePath: payload => this.record(
  216. 'session.openWorkspacePath',
  217. payload,
  218. this.onOpenWorkspacePath(payload),
  219. ),
  220. page: request => this.page(request),
  221. follow: (request, signal) => this.openFollow(request, signal),
  222. control: signal => this.openControl(signal),
  223. },
  224. subagents: {
  225. list: parentSessionId => this.record(
  226. 'subagents.list',
  227. parentSessionId,
  228. this.onSubagentList(parentSessionId),
  229. ),
  230. prompt: request => this.record('subagents.prompt', request, this.onSubagentPrompt(request)),
  231. interruptByParent: (childSessionId, parentSessionId, mode) => this.record(
  232. 'subagents.interruptByParent',
  233. { childSessionId, parentSessionId, mode },
  234. this.onSubagentInterrupt({ childSessionId, parentSessionId, mode }),
  235. ),
  236. },
  237. workspace: {
  238. create: payload => this.record('workspace.create', payload, this.onWorkspaceCreate(payload)),
  239. rename: payload => this.record('workspace.rename', payload, this.onWorkspaceRename(payload)),
  240. delete: payload => this.record('workspace.delete', payload, this.onWorkspaceDelete(payload)),
  241. insertBefore: payload => this.record(
  242. 'workspace.insertBefore',
  243. payload,
  244. this.onWorkspaceInsertBefore(payload),
  245. ),
  246. insertSessionBefore: payload => this.record(
  247. 'workspace.insertSessionBefore',
  248. payload,
  249. this.onWorkspaceInsertSessionBefore(payload),
  250. ),
  251. archiveSession: payload => this.record(
  252. 'workspace.archiveSession',
  253. payload,
  254. this.onWorkspaceArchiveSession(payload),
  255. ),
  256. follow: signal => this.openWorkspace(signal),
  257. },
  258. }
  259. }
  260. /** Push one live Session event to every follower of that Session. */
  261. async pushFollow(
  262. sessionId: SessionId,
  263. frame: Exclude<SessionFollowFrame, { type: 'snapshot' }>,
  264. ): Promise<void> {
  265. await Promise.all([...(this.followConns.get(sessionId) ?? [])].map(conn => new Promise<void>((resolve) => {
  266. conn.feed({ kind: 'frame', value: frame, delivered: resolve })
  267. })))
  268. }
  269. /** Push one Host-wide control update. */
  270. pushControl(frame: Exclude<SessionControlFrame, { type: 'baseline' }>): void {
  271. for (const conn of [...this.controlConns]) conn.feed({ kind: 'frame', value: frame })
  272. }
  273. /** Push one Workspace projection increment. */
  274. pushWorkspace(frame: Exclude<WorkspaceFollowFrame, { type: 'baseline' }>): void {
  275. for (const conn of [...this.workspaceConns]) conn.feed({ kind: 'frame', value: frame })
  276. }
  277. /** End (clean close) or fail (throw) every open stream — reconnect-path material. */
  278. endStreams(): void {
  279. for (const conns of this.followConns.values()) {
  280. for (const conn of [...conns]) conn.feed({ kind: 'end' })
  281. }
  282. for (const conn of [...this.controlConns]) conn.feed({ kind: 'end' })
  283. for (const conn of [...this.workspaceConns]) conn.feed({ kind: 'end' })
  284. }
  285. failStreams(error: unknown): void {
  286. for (const conns of this.followConns.values()) {
  287. for (const conn of [...conns]) conn.feed({ kind: 'fail', error })
  288. }
  289. for (const conn of [...this.controlConns]) conn.feed({ kind: 'fail', error })
  290. for (const conn of [...this.workspaceConns]) conn.feed({ kind: 'fail', error })
  291. }
  292. callsOf(method: string): unknown[] {
  293. return this.calls.filter(c => c.method === method).map(c => c.payload)
  294. }
  295. /** Number of currently attached journal generations for one Session. */
  296. activeFollows(sessionId: SessionId): number {
  297. return this.followConns.get(sessionId)?.length ?? 0
  298. }
  299. private record<T>(method: string, payload: unknown, response: Promise<T>): Promise<T> {
  300. this.calls.push({ method, payload })
  301. return response
  302. }
  303. private page(request: SessionPageRequest): Promise<RemoteResult<SessionPage>> {
  304. return this.fetchPage(request)
  305. }
  306. private async fetchPage(
  307. request: SessionPageRequest,
  308. response?: Promise<RemoteResult<SessionPage>>,
  309. ): Promise<RemoteResult<SessionPage>> {
  310. const sessionId = addressSessionId(request.address)
  311. const payload = request.address.kind === 'session'
  312. ? {
  313. sessionId,
  314. throughSeq: request.throughSeq,
  315. ...request.beforeSeq === undefined ? {} : { beforeSeq: request.beforeSeq },
  316. ...request.maxMessages === undefined ? {} : { maxMessages: request.maxMessages },
  317. }
  318. : {
  319. parentSessionId: request.address.parentSessionId,
  320. childSessionId: request.address.childSessionId,
  321. mode: request.address.mode,
  322. throughSeq: request.throughSeq,
  323. ...request.beforeSeq === undefined ? {} : { beforeSeq: request.beforeSeq },
  324. ...request.maxMessages === undefined ? {} : { maxMessages: request.maxMessages },
  325. }
  326. const method = request.address.kind === 'session' ? 'session.history' : 'subagent.history'
  327. const result = await this.record(method, payload, response ?? this.onHistory({
  328. sessionId,
  329. throughSeq: request.throughSeq,
  330. ...request.beforeSeq === undefined ? {} : { beforeSeq: request.beforeSeq },
  331. ...request.maxMessages === undefined ? {} : { maxMessages: request.maxMessages },
  332. }))
  333. if (!result.ok) return result
  334. return {
  335. ok: true,
  336. value: {
  337. ...result.value,
  338. records: result.value.records
  339. .filter(record => historyRecordLastSeq(record) <= request.throughSeq),
  340. },
  341. }
  342. }
  343. private async *openFollow(
  344. request: SessionFollowRequest,
  345. signal: AbortSignal = new AbortController().signal,
  346. ): AsyncGenerator<SessionFollowFrame> {
  347. const sessionId = addressSessionId(request.address)
  348. this.followStarts.push(sessionId)
  349. this.calls.push({ method: 'session.follow', payload: request })
  350. const conns = this.followConns.get(sessionId) ?? []
  351. if (!this.followConns.has(sessionId)) this.followConns.set(sessionId, conns)
  352. const stream = this.openValueStream(conns, signal)
  353. try {
  354. const response = await this.onHistory({
  355. sessionId,
  356. maxMessages: request.maxMessages ?? 50,
  357. })
  358. if (!response.ok) throw response.error
  359. const page = response.value
  360. const tail = page.records.at(-1)
  361. const cursor = this.followCursor ?? (tail === undefined ? -1 : historyRecordLastSeq(tail))
  362. yield {
  363. type: 'snapshot',
  364. header: {
  365. version: 1,
  366. id: sessionId,
  367. createdAt: 0,
  368. ...(request.address.kind === 'subagent'
  369. ? { origin: 'subagent' as const, parentSession: request.address.parentSessionId }
  370. : {}),
  371. },
  372. cursor,
  373. records: page.records.filter(record => historyRecordLastSeq(record) <= cursor),
  374. hasMore: page.hasMore,
  375. projections: page.projections ?? { asOfSeq: cursor, values: {} },
  376. ...request.assistantStream === true
  377. ? { assistantStream: this.assistantStreamBaseline }
  378. : {},
  379. }
  380. yield* stream.values
  381. } finally {
  382. stream.dispose()
  383. }
  384. }
  385. private async *openControl(
  386. signal: AbortSignal = new AbortController().signal,
  387. ): AsyncGenerator<SessionControlFrame> {
  388. const stream = this.openValueStream(this.controlConns, signal)
  389. try {
  390. yield { type: 'baseline', value: this.controlBaseline }
  391. yield* stream.values
  392. } finally {
  393. stream.dispose()
  394. }
  395. }
  396. private async *openWorkspace(
  397. signal: AbortSignal = new AbortController().signal,
  398. ): AsyncGenerator<WorkspaceFollowFrame> {
  399. const stream = this.openValueStream(this.workspaceConns, signal)
  400. try {
  401. yield { type: 'baseline', value: this.workspaceBaseline }
  402. yield* stream.values
  403. } finally {
  404. stream.dispose()
  405. }
  406. }
  407. private openValueStream<F>(
  408. registry: ValueStreamConn<F>[],
  409. signal: AbortSignal,
  410. ): OpenValueStream<F> {
  411. const inbox: ValueStreamItem<F>[] = []
  412. let wake: (() => void) | null = null
  413. let inFlightDelivered: (() => void) | undefined
  414. let disposed = false
  415. const conn: ValueStreamConn<F> = {
  416. feed: (item) => {
  417. inbox.push(item)
  418. wake?.()
  419. },
  420. }
  421. registry.push(conn)
  422. const dispose = (): void => {
  423. if (disposed) return
  424. disposed = true
  425. inFlightDelivered?.()
  426. for (const item of inbox) {
  427. if (item.kind === 'frame') item.delivered?.()
  428. }
  429. const index = registry.indexOf(conn)
  430. if (index >= 0) registry.splice(index, 1)
  431. wake?.()
  432. }
  433. const values = (async function* (): AsyncGenerator<F> {
  434. try {
  435. while (!signal.aborted && !disposed) {
  436. while (inbox.length > 0) {
  437. const item = inbox.shift() as ValueStreamItem<F>
  438. if (item.kind === 'end') return
  439. if (item.kind === 'fail') throw item.error
  440. inFlightDelivered = item.delivered
  441. yield item.value
  442. inFlightDelivered?.()
  443. inFlightDelivered = undefined
  444. }
  445. await new Promise<void>((resolve) => {
  446. wake = resolve
  447. signal.addEventListener('abort', () => { resolve() }, { once: true })
  448. })
  449. wake = null
  450. }
  451. } finally {
  452. dispose()
  453. }
  454. })()
  455. return { values, dispose }
  456. }
  457. }