fake-api.client.ts 21 KB

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