fake-api.client.ts 24 KB

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