fake-api.client.ts 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550
  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). Session streams are hand pumps: pushFollow/pushControl.
  4. import type {
  5. IApiClient,
  6. RpcError, RpcResponse, SessionId, SessionSearchItem, SkillEntry,
  7. WorkspaceId, WorkspaceView,
  8. } from '@deepseek-ai/dsh-api-remotes/client'
  9. import type {
  10. SessionAddress,
  11. SessionControlBaseline,
  12. SessionControlFrame,
  13. SessionFollowFrame,
  14. SessionFollowRequest,
  15. SessionPage,
  16. SessionPageRequest,
  17. SessionProjectionBaseline,
  18. SessionSelectModelRequest,
  19. SessionSelectModelValue,
  20. } from '@deepseek-ai/dsh-api-session-controller/types'
  21. import type { WorkspaceRemote } from '@deepseek-ai/dsh-api-workspace-controller/client'
  22. import type { WorkspaceFollowFrame } from '@deepseek-ai/dsh-api-workspace-controller/types'
  23. import type { RemoteResult } from '@deepseek-ai/dsh-typert-protocol'
  24. import {
  25. RemoteStream,
  26. RemoteStreamError,
  27. type RemoteStreamOptions,
  28. } from '@deepseek-ai/dsh-api-gateway/client'
  29. import { RpcId } from '@deepseek-ai/dsh-client-connection/client'
  30. import type { SessionRemotes } from '../src/client/sessions/remotes.ts'
  31. import { historyRecordLastSeq } from '../src/client/sessions/history-records.ts'
  32. const AVAILABLE_STREAM_CONNECTION = {
  33. hostDescription: {
  34. getSnapshot: () => ({
  35. version: 'fixture', cwd: '/f', attachedSessions: 0, home: '/h', canOpenPath: true,
  36. }),
  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. let nextRpc = 0
  71. export function ok<T>(value: T): RpcResponse<T> {
  72. return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: true, value } }
  73. }
  74. export function err<T>(error: RpcError): RpcResponse<T> {
  75. return { rpcId: RpcId(`fake-${nextRpc++}`), result: { ok: false, error } }
  76. }
  77. /** Successful generated Remote result for programmable domain fakes. */
  78. function remoteOk<T>(value: T): RemoteResult<T> {
  79. return { ok: true, value }
  80. }
  81. type ValueStreamItem<F> =
  82. | { kind: 'frame'; value: F; delivered?: () => void }
  83. | { kind: 'end' }
  84. | { kind: 'fail'; error: unknown }
  85. interface ValueStreamConn<F> {
  86. feed(item: ValueStreamItem<F>): void
  87. }
  88. interface OpenValueStream<F> {
  89. readonly values: AsyncGenerator<F>
  90. dispose(): void
  91. }
  92. /**
  93. * Commands Remote double: the generated face delivers the carrier's outcome, so
  94. * a test that programs nothing sees an empty catalog and an unmatched line.
  95. * @returns the Remote namespaces the session cluster calls.
  96. */
  97. export type RuntimeRemotes = SessionRemotes & { readonly workspace: WorkspaceRemote }
  98. export function fakeRemote(api = new FakeApiClient()): RuntimeRemotes {
  99. return api.sessionRemotes()
  100. }
  101. export class FakeApiClient implements IApiClient {
  102. /** Chronological call record: [method, payload]. */
  103. readonly calls: { method: string; payload: unknown }[] = []
  104. /** Session ids in physical follow-generation opening order. */
  105. readonly followStarts: SessionId[] = []
  106. // Programmable slots (defaults answer OK-empty); reassign per case.
  107. onList: (payload: unknown) => Promise<RpcResponse<{ items: never[] }>> = () => Promise.resolve(ok({ items: [] }))
  108. onSearch: (payload: unknown) => Promise<RpcResponse<{ items: SessionSearchItem[]; hasMore: boolean }>> =
  109. () => Promise.resolve(ok({ items: [], hasMore: false }))
  110. onCreate: (payload: unknown) => Promise<RpcResponse<{ sessionId: SessionId }>> = () => Promise.resolve(ok({ sessionId: 'fk-new' as SessionId }))
  111. onSelectModel: (payload: SessionSelectModelRequest) => Promise<RpcResponse<SessionSelectModelValue>> =
  112. payload => Promise.resolve(ok({
  113. selected: {
  114. provider: payload.provider,
  115. model: payload.model,
  116. ...(payload.reasoningEffort === undefined
  117. ? {}
  118. : { reasoningEffort: payload.reasoningEffort }),
  119. },
  120. }))
  121. onRename: (payload: unknown) => Promise<RpcResponse<{ title: string; seq: number }>> = () => Promise.resolve(ok({ title: 'fk-renamed', seq: 0 }))
  122. onFork: (payload: unknown) => Promise<RpcResponse<{ sessionId: SessionId }>> = () => Promise.resolve(ok({ sessionId: 'fk-fork' as SessionId }))
  123. onHistory: (payload: { sessionId: SessionId; throughSeq?: number; beforeSeq?: number; maxMessages?: number })
  124. => Promise<RpcResponse<SessionPage & { readonly projections?: SessionProjectionBaseline }>> =
  125. () => Promise.resolve(ok({ records: [], hasMore: false }))
  126. onPrompt: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  127. onAttachment: (payload: unknown) => Promise<RpcResponse<{ attachment: { attachmentId: never; mediaType: 'image/png'; bytes: number; width: number; height: number }; data: string }>> =
  128. () => Promise.resolve(ok({ attachment: { attachmentId: 'a' as never, mediaType: 'image/png', bytes: 1, width: 1, height: 1 }, data: 'AA==' }))
  129. onUpdateQueue: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  130. onCancel: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>> = () => Promise.resolve(ok({ accepted: true as const }))
  131. onDescribe: (payload: unknown) => Promise<RpcResponse<{
  132. version: string
  133. cwd: string
  134. attachedSessions: number
  135. home: string
  136. canOpenPath: boolean
  137. }>> =
  138. () => Promise.resolve(ok({
  139. version: '0-fake', cwd: '/f', attachedSessions: 0, home: '/h', canOpenPath: true,
  140. }))
  141. onPickDirectory: (payload: unknown) => Promise<RpcResponse<{ path: string | null }>> =
  142. () => Promise.resolve(ok({ path: null }))
  143. onOpenPath: (payload: unknown) => Promise<RpcResponse<{ opened: true }>> =
  144. () => Promise.resolve(ok({ opened: true as const }))
  145. onListDirectory: (payload: unknown) => Promise<RpcResponse<{
  146. path: string
  147. home: string
  148. crumbs: { name: string; path: string; hidden: boolean }[]
  149. entries: { name: string; path: string; hidden: boolean }[]
  150. truncated: boolean
  151. }>> =
  152. () => Promise.resolve(ok({ path: '/home/fake', home: '/home/fake', crumbs: [{ name: '/', path: '/', hidden: false }], entries: [], truncated: false }))
  153. onCreateDirectory: (payload: unknown) => Promise<RpcResponse<{ path: string }>> =
  154. () => Promise.resolve(ok({ path: '/home/fake/new' }))
  155. private readonly followConns = new Map<SessionId, ValueStreamConn<SessionFollowFrame>[]>()
  156. private readonly controlConns: ValueStreamConn<SessionControlFrame>[] = []
  157. private readonly workspaceConns: ValueStreamConn<WorkspaceFollowFrame>[] = []
  158. /** Optional Host opening cursor override for stale-page and reconnect tests. */
  159. followCursor: number | undefined
  160. controlBaseline: SessionControlBaseline = {
  161. queues: {},
  162. jobs: {},
  163. projections: {},
  164. }
  165. workspaceBaseline: Extract<WorkspaceFollowFrame, { type: 'baseline' }>['value'] = {
  166. items: [],
  167. archivedSessionIds: [],
  168. }
  169. lastSearchSignal: AbortSignal | undefined
  170. onSubagentList: (payload: unknown) => Promise<RpcResponse<{ entries: never[]; parentAvailable: boolean }>>
  171. = () => Promise.resolve(ok({ entries: [], parentAvailable: true }))
  172. onSubagentPrompt: (payload: unknown) => Promise<RpcResponse<{ messageId: never }>>
  173. = () => Promise.resolve(ok({ messageId: 'fake-message' as never }))
  174. onSubagentInterrupt: (payload: unknown) => Promise<RpcResponse<{ accepted: true }>>
  175. = () => Promise.resolve(ok({ accepted: true as const }))
  176. readonly subagents: IApiClient['subagents'] = {
  177. list: (payload: unknown) => this.record('subagent.list', payload, this.onSubagentList(payload)),
  178. prompt: (payload: unknown) => this.record('subagent.prompt', payload, this.onSubagentPrompt(payload)),
  179. interrupt: (payload: unknown) => this.record('subagent.interrupt', payload, this.onSubagentInterrupt(payload)),
  180. }
  181. readonly host: IApiClient['host'] = {
  182. describe: (payload: unknown) => this.record('host.describe', payload, this.onDescribe(payload)),
  183. pickDirectory: (payload: unknown) => this.record('host.pickDirectory', payload, this.onPickDirectory(payload)),
  184. listDirectory: (payload: unknown) => this.record('host.listDirectory', payload, this.onListDirectory(payload)),
  185. createDirectory: (payload: unknown) => this.record('host.createDirectory', payload, this.onCreateDirectory(payload)),
  186. openPath: (payload: unknown) => this.record('host.openPath', payload, this.onOpenPath(payload)),
  187. }
  188. onWorkspaceCreate: (payload: unknown) => Promise<RemoteResult<{ workspace: WorkspaceView; created: boolean }>> =
  189. () => Promise.resolve(remoteOk({ workspace: fakeWorkspace('fk-ws'), created: true }))
  190. onWorkspaceRename: (payload: unknown) => Promise<RemoteResult<{ workspace: WorkspaceView }>> =
  191. () => Promise.resolve(remoteOk({ workspace: fakeWorkspace('fk-ws') }))
  192. onWorkspaceDelete: (payload: unknown) => Promise<RemoteResult<{ deleted: true }>> =
  193. () => Promise.resolve(remoteOk({ deleted: true }))
  194. onWorkspaceInsertBefore: (payload: unknown) => Promise<RemoteResult<{ workspaceIds: WorkspaceId[] }>> =
  195. () => Promise.resolve(remoteOk({ workspaceIds: [] }))
  196. onWorkspaceInsertSessionBefore: (payload: unknown) => Promise<RemoteResult<{ workspace: WorkspaceView }>> =
  197. () => Promise.resolve(remoteOk({ workspace: fakeWorkspace('fk-ws') }))
  198. onWorkspaceArchiveSession: (payload: unknown) => Promise<RemoteResult<{ archivedSessionIds: SessionId[] }>> =
  199. payload => Promise.resolve(remoteOk({ archivedSessionIds: [(payload as { sessionId: SessionId }).sessionId] }))
  200. // Payloads stay `unknown` (lint-lane note above); response rows are the real
  201. // wire shapes so cases can program requires-bearing catalogs and dual-address
  202. // skill lists without casts.
  203. onSkillList: (payload: unknown) => Promise<RpcResponse<{ skills: SkillEntry[] }>>
  204. = () => Promise.resolve(ok({ skills: [] }))
  205. readonly agentPresets: IApiClient['agentPresets'] = {
  206. openDocument: (payload: { agentPreset: string }) =>
  207. this.record('agentPreset.openDocument', payload, Promise.resolve(ok({ opened: true as const }))),
  208. }
  209. readonly skills: IApiClient['skills'] = {
  210. list: (payload: unknown) => this.record('skill.list', payload, this.onSkillList(payload)),
  211. }
  212. readonly settings: IApiClient['settings'] = {
  213. describe: payload => this.record('settings.describe', payload, Promise.resolve(ok({ writable: true, hasDocument: false, namespaces: [] }))),
  214. openDocument: payload => this.record('settings.openDocument', payload, Promise.resolve(ok({ opened: true as const }))),
  215. update: payload => this.record('settings.update', payload, Promise.resolve(ok({ ns: 'fake', schema: {}, value: {}, applies: 'live' as const, secrets: [], revision: 0 }))),
  216. replace: payload => this.record('settings.replace', payload, Promise.resolve(ok({ ns: 'fake', schema: {}, value: {}, applies: 'live' as const, secrets: [], revision: 0 }))),
  217. mutate: payload => this.record('settings.mutate', payload, Promise.resolve(ok({ ns: 'fake', schema: {}, value: {}, applies: 'live' as const, secrets: [], revision: 0 }))),
  218. }
  219. readonly credentials: IApiClient['credentials'] = {
  220. describe: payload => this.record('credentials.describe', payload, Promise.resolve(ok({ credentials: {} }))),
  221. set: payload => this.record('credentials.set', payload, Promise.resolve(ok({}))),
  222. unset: payload => this.record('credentials.unset', payload, Promise.resolve(ok({}))),
  223. }
  224. readonly llm: IApiClient['llm'] = {
  225. providers: payload => this.record('llm.providers', payload, Promise.resolve(ok({ providers: [] }))),
  226. models: payload => this.record('llm.models', payload, Promise.resolve(ok({
  227. default: { provider: 'fixture', model: 'fixture' },
  228. routableProviders: [],
  229. groups: [],
  230. failures: [],
  231. }))),
  232. discoverModels: payload => this.record('llm.discoverModels', payload, Promise.resolve(ok({ models: [] }))),
  233. }
  234. /** Remote namespaces bound to this fake's programmable unary slots and stream pumps. */
  235. sessionRemotes(): RuntimeRemotes {
  236. return {
  237. $stream: <Item>(options: RemoteStreamOptions<Item>) => (
  238. new RemoteStream(AVAILABLE_STREAM_CONNECTION, options)
  239. ),
  240. commands: {
  241. execute: () => Promise.resolve({ ok: true, value: undefined }),
  242. },
  243. session: {
  244. list: payload => this.remoteResult('session.list', payload, this.onList(payload)),
  245. search: (payload, signal) => {
  246. this.lastSearchSignal = signal
  247. return this.remoteResult('session.search', payload, this.onSearch(payload))
  248. },
  249. create: payload => this.remoteResult('session.create', payload, this.onCreate(payload)),
  250. selectModel: payload => this.remoteResult(
  251. 'session.selectModel',
  252. payload,
  253. this.onSelectModel(payload),
  254. ),
  255. rename: payload => this.remoteResult('session.rename', payload, this.onRename(payload)),
  256. fork: payload => this.remoteResult('session.fork', payload, this.onFork(payload)),
  257. prompt: payload => this.remoteResult('session.prompt', payload, this.onPrompt(payload)),
  258. attachment: payload => this.remoteResult('session.attachment', payload, this.onAttachment(payload)),
  259. updateQueue: payload => this.remoteResult('session.updateQueue', payload, this.onUpdateQueue(payload)),
  260. cancel: payload => this.remoteResult('session.cancel', payload, this.onCancel(payload)),
  261. page: request => this.page(request),
  262. follow: (request, signal) => this.openFollow(request, signal),
  263. control: signal => this.openControl(signal),
  264. },
  265. workspace: {
  266. create: payload => this.record('workspace.create', payload, this.onWorkspaceCreate(payload)),
  267. rename: payload => this.record('workspace.rename', payload, this.onWorkspaceRename(payload)),
  268. delete: payload => this.record('workspace.delete', payload, this.onWorkspaceDelete(payload)),
  269. insertBefore: payload => this.record(
  270. 'workspace.insertBefore',
  271. payload,
  272. this.onWorkspaceInsertBefore(payload),
  273. ),
  274. insertSessionBefore: payload => this.record(
  275. 'workspace.insertSessionBefore',
  276. payload,
  277. this.onWorkspaceInsertSessionBefore(payload),
  278. ),
  279. archiveSession: payload => this.record(
  280. 'workspace.archiveSession',
  281. payload,
  282. this.onWorkspaceArchiveSession(payload),
  283. ),
  284. follow: signal => this.openWorkspace(signal),
  285. },
  286. }
  287. }
  288. /** Push one live Session event to every follower of that Session. */
  289. async pushFollow(
  290. sessionId: SessionId,
  291. frame: Extract<SessionFollowFrame, { type: 'event' }>,
  292. ): Promise<void> {
  293. await Promise.all([...(this.followConns.get(sessionId) ?? [])].map(conn => new Promise<void>((resolve) => {
  294. conn.feed({ kind: 'frame', value: frame, delivered: resolve })
  295. })))
  296. }
  297. /** Push one Host-wide control update. */
  298. pushControl(frame: Exclude<SessionControlFrame, { type: 'baseline' }>): void {
  299. for (const conn of [...this.controlConns]) conn.feed({ kind: 'frame', value: frame })
  300. }
  301. /** Push one Workspace projection increment. */
  302. pushWorkspace(frame: Exclude<WorkspaceFollowFrame, { type: 'baseline' }>): void {
  303. for (const conn of [...this.workspaceConns]) conn.feed({ kind: 'frame', value: frame })
  304. }
  305. /** End (clean close) or fail (throw) every open stream — reconnect-path material. */
  306. endStreams(): void {
  307. for (const conns of this.followConns.values()) {
  308. for (const conn of [...conns]) conn.feed({ kind: 'end' })
  309. }
  310. for (const conn of [...this.controlConns]) conn.feed({ kind: 'end' })
  311. for (const conn of [...this.workspaceConns]) conn.feed({ kind: 'end' })
  312. }
  313. failStreams(error: unknown): void {
  314. for (const conns of this.followConns.values()) {
  315. for (const conn of [...conns]) conn.feed({ kind: 'fail', error })
  316. }
  317. for (const conn of [...this.controlConns]) conn.feed({ kind: 'fail', error })
  318. for (const conn of [...this.workspaceConns]) conn.feed({ kind: 'fail', error })
  319. }
  320. callsOf(method: string): unknown[] {
  321. return this.calls.filter(c => c.method === method).map(c => c.payload)
  322. }
  323. /** Number of currently attached journal generations for one Session. */
  324. activeFollows(sessionId: SessionId): number {
  325. return this.followConns.get(sessionId)?.length ?? 0
  326. }
  327. private record<T>(method: string, payload: unknown, response: Promise<T>): Promise<T> {
  328. this.calls.push({ method, payload })
  329. return response
  330. }
  331. private async remoteResult<T>(
  332. method: string,
  333. payload: unknown,
  334. response: Promise<RpcResponse<T>>,
  335. ): Promise<RemoteResult<T>> {
  336. return (await this.record(method, payload, response)).result
  337. }
  338. private page(request: SessionPageRequest): Promise<RemoteResult<SessionPage>> {
  339. return this.fetchPage(request)
  340. }
  341. private async fetchPage(
  342. request: SessionPageRequest,
  343. response?: Promise<RpcResponse<SessionPage>>,
  344. ): Promise<RemoteResult<SessionPage>> {
  345. const sessionId = addressSessionId(request.address)
  346. const payload = request.address.kind === 'session'
  347. ? {
  348. sessionId,
  349. throughSeq: request.throughSeq,
  350. ...request.beforeSeq === undefined ? {} : { beforeSeq: request.beforeSeq },
  351. ...request.maxMessages === undefined ? {} : { maxMessages: request.maxMessages },
  352. }
  353. : {
  354. parentSessionId: request.address.parentSessionId,
  355. childSessionId: request.address.childSessionId,
  356. mode: request.address.mode,
  357. throughSeq: request.throughSeq,
  358. ...request.beforeSeq === undefined ? {} : { beforeSeq: request.beforeSeq },
  359. ...request.maxMessages === undefined ? {} : { maxMessages: request.maxMessages },
  360. }
  361. const method = request.address.kind === 'session' ? 'session.history' : 'subagent.history'
  362. const result = await this.remoteResult(method, payload, response ?? this.onHistory({
  363. sessionId,
  364. throughSeq: request.throughSeq,
  365. ...request.beforeSeq === undefined ? {} : { beforeSeq: request.beforeSeq },
  366. ...request.maxMessages === undefined ? {} : { maxMessages: request.maxMessages },
  367. }))
  368. if (!result.ok) return result
  369. return {
  370. ok: true,
  371. value: {
  372. ...result.value,
  373. records: result.value.records
  374. .filter(record => historyRecordLastSeq(record) <= request.throughSeq),
  375. },
  376. }
  377. }
  378. private async *openFollow(
  379. request: SessionFollowRequest,
  380. signal: AbortSignal = new AbortController().signal,
  381. ): AsyncGenerator<SessionFollowFrame> {
  382. const sessionId = addressSessionId(request.address)
  383. this.followStarts.push(sessionId)
  384. this.calls.push({ method: 'session.follow', payload: request })
  385. const conns = this.followConns.get(sessionId) ?? []
  386. if (!this.followConns.has(sessionId)) this.followConns.set(sessionId, conns)
  387. const stream = this.openValueStream(conns, signal)
  388. try {
  389. const response = await this.onHistory({
  390. sessionId,
  391. maxMessages: request.maxMessages ?? 50,
  392. })
  393. if (!response.result.ok) {
  394. throw new RemoteStreamError(
  395. response.result.error.code,
  396. response.result.error.message,
  397. response.result.error.details,
  398. )
  399. }
  400. const page = response.result.value
  401. const tail = page.records.at(-1)
  402. const cursor = this.followCursor ?? (tail === undefined ? -1 : historyRecordLastSeq(tail))
  403. yield {
  404. type: 'snapshot',
  405. header: {
  406. version: 0,
  407. id: sessionId,
  408. createdAt: 0,
  409. ...(request.address.kind === 'subagent'
  410. ? { origin: 'subagent' as const, parentSession: request.address.parentSessionId }
  411. : {}),
  412. },
  413. cursor,
  414. records: page.records.filter(record => historyRecordLastSeq(record) <= cursor),
  415. hasMore: page.hasMore,
  416. projections: page.projections ?? { asOfSeq: cursor, values: {} },
  417. }
  418. yield* stream.values
  419. } finally {
  420. stream.dispose()
  421. }
  422. }
  423. private async *openControl(
  424. signal: AbortSignal = new AbortController().signal,
  425. ): AsyncGenerator<SessionControlFrame> {
  426. const stream = this.openValueStream(this.controlConns, signal)
  427. try {
  428. yield { type: 'baseline', value: this.controlBaseline }
  429. yield* stream.values
  430. } finally {
  431. stream.dispose()
  432. }
  433. }
  434. private async *openWorkspace(
  435. signal: AbortSignal = new AbortController().signal,
  436. ): AsyncGenerator<WorkspaceFollowFrame> {
  437. const stream = this.openValueStream(this.workspaceConns, signal)
  438. try {
  439. yield { type: 'baseline', value: this.workspaceBaseline }
  440. yield* stream.values
  441. } finally {
  442. stream.dispose()
  443. }
  444. }
  445. private openValueStream<F>(
  446. registry: ValueStreamConn<F>[],
  447. signal: AbortSignal,
  448. ): OpenValueStream<F> {
  449. const inbox: ValueStreamItem<F>[] = []
  450. let wake: (() => void) | null = null
  451. let inFlightDelivered: (() => void) | undefined
  452. let disposed = false
  453. const conn: ValueStreamConn<F> = {
  454. feed: (item) => {
  455. inbox.push(item)
  456. wake?.()
  457. },
  458. }
  459. registry.push(conn)
  460. const dispose = (): void => {
  461. if (disposed) return
  462. disposed = true
  463. inFlightDelivered?.()
  464. for (const item of inbox) {
  465. if (item.kind === 'frame') item.delivered?.()
  466. }
  467. const index = registry.indexOf(conn)
  468. if (index >= 0) registry.splice(index, 1)
  469. wake?.()
  470. }
  471. const values = (async function* (): AsyncGenerator<F> {
  472. try {
  473. while (!signal.aborted && !disposed) {
  474. while (inbox.length > 0) {
  475. const item = inbox.shift() as ValueStreamItem<F>
  476. if (item.kind === 'end') return
  477. if (item.kind === 'fail') throw item.error
  478. inFlightDelivered = item.delivered
  479. yield item.value
  480. inFlightDelivered?.()
  481. inFlightDelivered = undefined
  482. }
  483. await new Promise<void>((resolve) => {
  484. wake = resolve
  485. signal.addEventListener('abort', () => { resolve() }, { once: true })
  486. })
  487. wake = null
  488. }
  489. } finally {
  490. dispose()
  491. }
  492. })()
  493. return { values, dispose }
  494. }
  495. }