fixture.client.spec.ts 79 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782
  1. import { afterEach, describe, expect, expectTypeOf, it, vi } from 'vitest'
  2. import type {
  3. RpcRequest,
  4. RpcResponse,
  5. RpcResult,
  6. SessionEvent,
  7. SessionId,
  8. StreamChunk,
  9. } from '../src/client/api.ts'
  10. import { RpcId } from '../src/client/api.ts'
  11. import {
  12. createFixtureConnectionRpc,
  13. createFixtureFaces,
  14. type FixtureAssistantStreamFrame,
  15. type FixtureOptions,
  16. } from '../src/client/fixture.ts'
  17. import type {
  18. ClientConnectionRpc, ConnectionRpcResult,
  19. } from '../src/rpc.ts'
  20. import type { DirectoryListing } from '@deepseek-ai/dsh-host-directory-picker/types'
  21. import type {
  22. ModelCatalog, ModelSelection, SessionAssistantStreamFrame,
  23. } from '@deepseek-ai/dsh-api-session-controller/types'
  24. const sid = (id: string): SessionId => id as SessionId
  25. type WorkspaceId = string & { readonly __fixtureWorkspaceId: 'WorkspaceId' }
  26. const req = <P>(payload: P): RpcRequest<P> => ({ rpcId: RpcId(`t-${Math.abs(Math.sin(reqCount++)).toString(36).slice(2, 10)}`), payload })
  27. let reqCount = 0
  28. interface FixtureSessionSummary {
  29. sessionId: SessionId
  30. updatedAt: number
  31. running: boolean
  32. blank: boolean
  33. parentSessionId?: SessionId
  34. origin?: 'subagent'
  35. cwd?: string
  36. agentPreset?: string
  37. }
  38. interface FixtureHistoryEntry {
  39. readonly type: 'event'
  40. readonly event: SessionEvent
  41. }
  42. type FixtureHistoryRecord = FixtureHistoryEntry
  43. interface FixturePage {
  44. readonly records: readonly FixtureHistoryRecord[]
  45. readonly hasMore: boolean
  46. }
  47. function historyEvents(records: readonly FixtureHistoryRecord[]): SessionEvent[] {
  48. return records.map(record => record.event)
  49. }
  50. function isReasoningDeltaChunk(
  51. value: unknown,
  52. ): value is Extract<StreamChunk, { readonly type: 'reasoning-delta' }> {
  53. if (typeof value !== 'object' || value === null) return false
  54. const record = value as { readonly type?: unknown; readonly text?: unknown }
  55. return record.type === 'reasoning-delta' && typeof record.text === 'string'
  56. }
  57. type FixtureFollowFrame =
  58. | {
  59. readonly type: 'snapshot'
  60. readonly cursor: number
  61. readonly records: readonly FixtureHistoryRecord[]
  62. readonly hasMore: boolean
  63. readonly projections: {
  64. readonly asOfSeq: number
  65. readonly values: Readonly<Record<string, unknown>>
  66. }
  67. readonly assistantStream?: {
  68. readonly revision: number
  69. readonly activeAttempt?: {
  70. readonly attemptId: string
  71. readonly startedAfterSeq: number
  72. readonly turn: number
  73. readonly step: number
  74. readonly nextIndex: number
  75. readonly stream: readonly unknown[]
  76. }
  77. }
  78. }
  79. | FixtureHistoryEntry
  80. | { readonly type: 'assistant-stream'; readonly frame: FixtureAssistantStreamFrame }
  81. type FixtureControlFrame =
  82. | {
  83. readonly type: 'baseline'
  84. readonly value: {
  85. readonly queues: Readonly<Record<string, readonly unknown[]>>
  86. readonly jobs: Readonly<Record<string, readonly unknown[]>>
  87. readonly approvals: readonly unknown[]
  88. readonly questions: readonly unknown[]
  89. readonly projections: Readonly<Record<string, {
  90. readonly asOfSeq: number
  91. readonly values: Readonly<Record<string, unknown>>
  92. }>>
  93. }
  94. }
  95. | {
  96. readonly type: 'projection'
  97. readonly sessionId: SessionId
  98. readonly key: string
  99. readonly value: unknown
  100. readonly seq: number
  101. }
  102. interface FixtureSessionRequests {
  103. list: { readonly cursor?: string }
  104. search: { readonly query: string }
  105. create: {
  106. readonly workspaceId?: WorkspaceId
  107. readonly cwd?: string
  108. readonly sessionId?: SessionId
  109. readonly agentPreset?: string
  110. }
  111. history: {
  112. readonly sessionId: SessionId
  113. readonly beforeSeq?: number
  114. readonly maxMessages?: number
  115. }
  116. selectModel: {
  117. readonly sessionId: SessionId
  118. readonly provider: string
  119. readonly model: string
  120. readonly reasoningEffort?: string
  121. }
  122. prompt: {
  123. readonly sessionId: SessionId
  124. readonly mode: 'queue' | 'steer'
  125. readonly content: readonly ({ readonly type: 'text'; readonly text: string } | {
  126. readonly type: 'image'
  127. readonly mediaType: 'image/png' | 'image/jpeg' | 'image/webp' | 'image/gif'
  128. readonly data: string
  129. readonly name?: string
  130. })[]
  131. }
  132. cancel: { readonly sessionId: SessionId }
  133. rename: { readonly sessionId: SessionId; readonly title: string }
  134. }
  135. interface FixtureSessionValues {
  136. list: { readonly items: FixtureSessionSummary[] }
  137. search: { readonly items: readonly { readonly sessionId: SessionId; readonly snippet: string }[]; readonly hasMore: boolean }
  138. create: { readonly sessionId: SessionId }
  139. history: FixturePage
  140. selectModel: { readonly selected: ModelSelection }
  141. prompt: { readonly accepted: true }
  142. cancel: Record<never, never>
  143. rename: { readonly title: string; readonly seq: number }
  144. }
  145. type FixtureSessionApi = {
  146. [K in keyof FixtureSessionRequests]: (
  147. request: RpcRequest<FixtureSessionRequests[K]>,
  148. signal?: AbortSignal,
  149. ) => Promise<RpcResponse<FixtureSessionValues[K]>>
  150. }
  151. type FixtureSessionClient = {
  152. [K in keyof FixtureSessionRequests]: (
  153. request: FixtureSessionRequests[K],
  154. signal?: AbortSignal,
  155. ) => Promise<RpcResponse<FixtureSessionValues[K]>>
  156. }
  157. interface FixtureSessionRemote {
  158. modelCatalog(): Promise<ConnectionRpcResult<ModelCatalog>>
  159. follow(sessionId: SessionId, signal: AbortSignal): AsyncIterable<FixtureFollowFrame>
  160. control(signal: AbortSignal): AsyncIterable<FixtureControlFrame>
  161. }
  162. interface FixtureWorkspaceView {
  163. readonly workspaceId: WorkspaceId
  164. readonly path: string
  165. readonly title: string
  166. readonly sessionIds: readonly SessionId[]
  167. readonly createdAt: string
  168. readonly updatedAt: string
  169. }
  170. interface FixtureWorkspaceRequests {
  171. create: { readonly path: string }
  172. rename: { readonly workspaceId: WorkspaceId; readonly title: string }
  173. delete: { readonly workspaceId: WorkspaceId }
  174. insertBefore: { readonly workspaceId: WorkspaceId; readonly beforeWorkspaceId?: WorkspaceId }
  175. insertSessionBefore: {
  176. readonly workspaceId: WorkspaceId
  177. readonly sessionId: SessionId
  178. readonly beforeSessionId?: SessionId
  179. }
  180. archiveSession: { readonly sessionId: SessionId }
  181. }
  182. interface FixtureWorkspaceValues {
  183. create: { readonly workspace: FixtureWorkspaceView; readonly created: boolean }
  184. rename: { readonly workspace: FixtureWorkspaceView }
  185. delete: { readonly deleted: true }
  186. insertBefore: { readonly workspaceIds: readonly WorkspaceId[] }
  187. insertSessionBefore: { readonly workspace: FixtureWorkspaceView }
  188. archiveSession: { readonly archivedSessionIds: readonly SessionId[] }
  189. }
  190. type FixtureWorkspaceApi = {
  191. [K in keyof FixtureWorkspaceRequests]: (
  192. request: RpcRequest<FixtureWorkspaceRequests[K]>,
  193. signal?: AbortSignal,
  194. ) => Promise<RpcResponse<FixtureWorkspaceValues[K]>>
  195. }
  196. type FixtureWorkspaceClient = {
  197. [K in keyof FixtureWorkspaceRequests]: (
  198. request: FixtureWorkspaceRequests[K],
  199. signal?: AbortSignal,
  200. ) => Promise<RpcResponse<FixtureWorkspaceValues[K]>>
  201. }
  202. type FixtureWorkspaceFrame =
  203. | {
  204. readonly type: 'baseline'
  205. readonly value: {
  206. readonly items: readonly FixtureWorkspaceView[]
  207. readonly archivedSessionIds: readonly SessionId[]
  208. }
  209. }
  210. | { readonly type: 'upsert'; readonly workspace: FixtureWorkspaceView }
  211. | { readonly type: 'remove'; readonly workspaceId: WorkspaceId }
  212. | { readonly type: 'order'; readonly workspaceIds: readonly WorkspaceId[] }
  213. | { readonly type: 'archived'; readonly archivedSessionIds: readonly SessionId[] }
  214. interface FixtureWorkspaceRemote {
  215. follow(signal: AbortSignal): AsyncIterable<FixtureWorkspaceFrame>
  216. }
  217. interface FixtureRemoteEventNotificationFrame {
  218. readonly type: 'emit'
  219. readonly event: string
  220. readonly args: readonly unknown[]
  221. }
  222. interface FixtureRemoteEventRequestFrame {
  223. readonly type: 'waterfall'
  224. readonly event: string
  225. readonly eventId: string
  226. readonly agentId: SessionId
  227. readonly request: Readonly<Record<string, unknown>>
  228. }
  229. interface FixtureRemoteEventCancellationFrame {
  230. readonly type: 'cancel'
  231. readonly eventId: string
  232. }
  233. type FixtureRemoteEventFrame =
  234. | FixtureRemoteEventNotificationFrame
  235. | FixtureRemoteEventRequestFrame
  236. | FixtureRemoteEventCancellationFrame
  237. interface FixtureRemoteEventResult {
  238. readonly clientId: string
  239. readonly eventId: string
  240. readonly outcome:
  241. | { readonly kind: 'next' }
  242. | { readonly kind: 'result'; readonly value?: unknown }
  243. | {
  244. readonly kind: 'rejected'
  245. readonly error: {
  246. readonly name: string
  247. readonly message: string
  248. readonly code?: string
  249. readonly details?: unknown
  250. }
  251. }
  252. }
  253. interface FixtureRemoteEventStream extends AsyncIterable<FixtureRemoteEventFrame> {
  254. readonly clientId: Promise<string>
  255. }
  256. type FixtureTestApi = {
  257. /** The directory-picking Remote namespace as the fixture serves it. */
  258. readonly directoryPickerRemote: {
  259. pick: () => Promise<ConnectionRpcResult<string | null>>
  260. list: (path?: string) => Promise<ConnectionRpcResult<DirectoryListing>>
  261. createDirectory: (path: string, name: string) => Promise<ConnectionRpcResult<string>>
  262. }
  263. readonly sessions: FixtureSessionApi
  264. readonly sessionRemote: FixtureSessionRemote
  265. readonly workspace: FixtureWorkspaceApi
  266. readonly workspaceRemote: FixtureWorkspaceRemote
  267. readonly credentialRemote: FixtureCredentialRemote
  268. readonly settingsRemote: FixtureSettingsRemote
  269. readonly remoteEvents: (signal: AbortSignal) => FixtureRemoteEventStream
  270. readonly answerRemoteEvent: (result: FixtureRemoteEventResult) => Promise<unknown>
  271. }
  272. /** Keep existing fixture assertions compact while driving only the new Session Remote endpoints. */
  273. function createFixtureApi(options: FixtureOptions = {}): FixtureTestApi {
  274. const { rpc } = createFixtureFaces(options)
  275. return {
  276. directoryPickerRemote: {
  277. pick: () => rpc.call('/api', 'directoryPicker/pick', { args: {} }) as
  278. Promise<ConnectionRpcResult<string | null>>,
  279. list: (path?: string) => rpc.call('/api', 'directoryPicker/list', { args: { path } }) as
  280. Promise<ConnectionRpcResult<DirectoryListing>>,
  281. createDirectory: (path: string, name: string) =>
  282. rpc.call('/api', 'directoryPicker/createDirectory', { args: { path, name } }) as
  283. Promise<ConnectionRpcResult<string>>,
  284. },
  285. sessions: createSessionApi(rpc),
  286. sessionRemote: createSessionRemote(rpc),
  287. workspace: createWorkspaceApi(rpc),
  288. workspaceRemote: createWorkspaceRemote(rpc),
  289. credentialRemote: createCredentialRemote(rpc),
  290. settingsRemote: createSettingsRemote(rpc),
  291. remoteEvents: (signal: AbortSignal) => openFixtureRemoteEvents(rpc, signal),
  292. answerRemoteEvent: (result: FixtureRemoteEventResult) =>
  293. rpc.call('/api', '$events/result', { args: result }),
  294. }
  295. }
  296. /** The fixture's Credentials Remote endpoints over the shared RPC carrier. */
  297. interface FixtureCredentialRemote {
  298. describe(refs: readonly string[]): Promise<ConnectionRpcResult<unknown>>
  299. set(ref: string, value: string): Promise<ConnectionRpcResult<unknown>>
  300. unset(ref: string): Promise<ConnectionRpcResult<unknown>>
  301. }
  302. /** The settings Remote reads the fixture serves, addressed like the credential half. */
  303. interface FixtureSettingsRemote {
  304. describe(): Promise<ConnectionRpcResult<unknown>>
  305. update(ns: string, patch: unknown, expectedRevision?: number): Promise<ConnectionRpcResult<unknown>>
  306. replace(ns: string, section: unknown, expectedRevision?: number): Promise<ConnectionRpcResult<unknown>>
  307. }
  308. function createSettingsRemote(rpc: ClientConnectionRpc): FixtureSettingsRemote {
  309. return {
  310. describe: () => rpc.call('/api', 'settings/describe', { args: {} }),
  311. update: (ns, patch, expectedRevision) => rpc.call('/api', 'settings/update', {
  312. args: { ns, patch, expectedRevision },
  313. }),
  314. replace: (ns, section, expectedRevision) => rpc.call('/api', 'settings/replace', {
  315. args: { ns, section, expectedRevision },
  316. }),
  317. }
  318. }
  319. function createCredentialRemote(rpc: ClientConnectionRpc): FixtureCredentialRemote {
  320. return {
  321. describe: refs => rpc.call('/api', 'credentials/describe', { args: { refs } }),
  322. set: (ref, value) => rpc.call('/api', 'credentials/set', { args: { ref, value } }),
  323. unset: ref => rpc.call('/api', 'credentials/unset', { args: { ref } }),
  324. }
  325. }
  326. function openFixtureRemoteEvents(
  327. rpc: ClientConnectionRpc,
  328. signal: AbortSignal,
  329. ): FixtureRemoteEventStream {
  330. const ready = Promise.withResolvers<string>()
  331. const source = (async function* (): AsyncGenerator<FixtureRemoteEventFrame> {
  332. const stream = rpc.open?.('/api', '$events', { args: {} }, signal)
  333. if (stream === undefined) throw new Error('fixture forwarded-event stream is unavailable')
  334. let opened = false
  335. for await (const value of stream) {
  336. if (!opened) {
  337. expect(value).toMatchObject({ type: 'ready' })
  338. const clientId: unknown = Reflect.get(value as object, 'clientId')
  339. if (typeof clientId !== 'string') throw new Error('fixture forwarded-event stream omitted its Client id')
  340. ready.resolve(clientId)
  341. opened = true
  342. continue
  343. }
  344. yield value as FixtureRemoteEventFrame
  345. }
  346. })()
  347. return Object.assign(source, { clientId: ready.promise })
  348. }
  349. function createSessionApi(rpc: ClientConnectionRpc): FixtureSessionApi {
  350. const call = async <K extends keyof FixtureSessionRequests>(
  351. endpoint: K,
  352. request: RpcRequest<FixtureSessionRequests[K]>,
  353. signal?: AbortSignal,
  354. ): Promise<RpcResponse<FixtureSessionValues[K]>> => {
  355. const page = endpoint === 'history'
  356. ? request.payload as FixtureSessionRequests['history']
  357. : undefined
  358. const args = endpoint === 'list'
  359. ? { _request: request.payload }
  360. : endpoint === 'history'
  361. ? {
  362. request: {
  363. address: { kind: 'session', sessionId: page?.sessionId },
  364. ...page?.beforeSeq === undefined ? {} : { beforeSeq: page.beforeSeq },
  365. ...page?.maxMessages === undefined ? {} : { maxMessages: page.maxMessages },
  366. },
  367. }
  368. : { request: request.payload }
  369. const remoteEndpoint = endpoint === 'history' ? 'page' : endpoint
  370. const result = await rpc.call('/api', `session/${remoteEndpoint}`, { args }, signal)
  371. return {
  372. rpcId: request.rpcId,
  373. result: result as unknown as RpcResult<FixtureSessionValues[K]>,
  374. }
  375. }
  376. return {
  377. list: (request, signal) => call('list', request, signal),
  378. search: (request, signal) => call('search', request, signal),
  379. create: (request, signal) => call('create', request, signal),
  380. history: (request, signal) => call('history', request, signal),
  381. selectModel: (request, signal) => call('selectModel', request, signal),
  382. prompt: (request, signal) => call('prompt', request, signal),
  383. cancel: (request, signal) => call('cancel', request, signal),
  384. rename: (request, signal) => call('rename', request, signal),
  385. }
  386. }
  387. function createSessionClient(rpc: ClientConnectionRpc): FixtureSessionClient {
  388. const api = createSessionApi(rpc)
  389. return {
  390. list: (request, signal) => api.list(req(request), signal),
  391. search: (request, signal) => api.search(req(request), signal),
  392. create: (request, signal) => api.create(req(request), signal),
  393. history: (request, signal) => api.history(req(request), signal),
  394. selectModel: (request, signal) => api.selectModel(req(request), signal),
  395. prompt: (request, signal) => api.prompt(req(request), signal),
  396. cancel: (request, signal) => api.cancel(req(request), signal),
  397. rename: (request, signal) => api.rename(req(request), signal),
  398. }
  399. }
  400. function createSessionRemote(rpc: ClientConnectionRpc): FixtureSessionRemote {
  401. const open = <F>(endpoint: string, args: object, signal: AbortSignal): AsyncIterable<F> => {
  402. const stream = rpc.open?.('/api', endpoint, { args }, signal)
  403. if (stream === undefined) throw new Error(`fixture ${endpoint} stream is unavailable`)
  404. return stream as AsyncIterable<F>
  405. }
  406. return {
  407. modelCatalog: () => rpc.call('/api', 'session/modelCatalog', { args: {} }) as
  408. Promise<ConnectionRpcResult<ModelCatalog>>,
  409. follow: (sessionId, signal) => open<FixtureFollowFrame>('session/follow', {
  410. request: { address: { kind: 'session', sessionId }, assistantStream: true },
  411. }, signal),
  412. control: signal => open<FixtureControlFrame>('session/control', {}, signal),
  413. }
  414. }
  415. function createWorkspaceApi(rpc: ClientConnectionRpc): FixtureWorkspaceApi {
  416. const call = async <K extends keyof FixtureWorkspaceRequests>(
  417. endpoint: K,
  418. request: RpcRequest<FixtureWorkspaceRequests[K]>,
  419. signal?: AbortSignal,
  420. ): Promise<RpcResponse<FixtureWorkspaceValues[K]>> => {
  421. const result = await rpc.call('/api', `workspace/${endpoint}`, {
  422. args: { request: request.payload },
  423. }, signal)
  424. return {
  425. rpcId: request.rpcId,
  426. result: result as unknown as RpcResult<FixtureWorkspaceValues[K]>,
  427. }
  428. }
  429. return {
  430. create: (request, signal) => call('create', request, signal),
  431. rename: (request, signal) => call('rename', request, signal),
  432. delete: (request, signal) => call('delete', request, signal),
  433. insertBefore: (request, signal) => call('insertBefore', request, signal),
  434. insertSessionBefore: (request, signal) => call('insertSessionBefore', request, signal),
  435. archiveSession: (request, signal) => call('archiveSession', request, signal),
  436. }
  437. }
  438. function createWorkspaceClient(rpc: ClientConnectionRpc): FixtureWorkspaceClient {
  439. const api = createWorkspaceApi(rpc)
  440. return {
  441. create: (request, signal) => api.create(req(request), signal),
  442. rename: (request, signal) => api.rename(req(request), signal),
  443. delete: (request, signal) => api.delete(req(request), signal),
  444. insertBefore: (request, signal) => api.insertBefore(req(request), signal),
  445. insertSessionBefore: (request, signal) => api.insertSessionBefore(req(request), signal),
  446. archiveSession: (request, signal) => api.archiveSession(req(request), signal),
  447. }
  448. }
  449. function createWorkspaceRemote(rpc: ClientConnectionRpc): FixtureWorkspaceRemote {
  450. return {
  451. follow(signal) {
  452. const stream = rpc.open?.('/api', 'workspace/follow', { args: {} }, signal)
  453. if (stream === undefined) throw new Error('fixture workspace/follow stream is unavailable')
  454. return stream as AsyncIterable<FixtureWorkspaceFrame>
  455. },
  456. }
  457. }
  458. interface TimingHooks {
  459. setHistoryDelay(ms: number): void
  460. failNextHistory(): void
  461. appendUser(id: string, msg: string): void
  462. appendTitle(id: string, title: string): void
  463. disarmOnlyGoal(): void
  464. startReasoningChunkStorm(id: string, chunkCount: number, chunksPerInterval: number, intervalMs: number): string
  465. reasoningChunkStormState(): {
  466. sessionId: string
  467. chunkCount: number
  468. chunksPerInterval: number
  469. intervalMs: number
  470. emitted: number
  471. marker: string
  472. emitting: boolean
  473. } | null
  474. beginModelRetry(id: string): void
  475. scheduleModelRetry(id: string, retry?: number, delayMs?: number): void
  476. cancelModelRetryDuringBackoff(id: string, delayMs?: number): void
  477. completeModelRetry(id: string): void
  478. appendSilent(id: string, msg: string): void
  479. breakStreams(): void
  480. }
  481. const timing = (): TimingHooks => (globalThis as Record<string, unknown>).__fxTiming as TimingHooks
  482. /** Collect value-stream frames until the predicate or a soft cap; abort ends the stream. */
  483. async function collectValues<F>(stream: AsyncIterable<F>, abort: AbortController, done: (frames: F[]) => boolean): Promise<F[]> {
  484. const frames: F[] = []
  485. for await (const frame of stream) {
  486. frames.push(frame)
  487. if (done(frames) || frames.length > 500) {
  488. abort.abort()
  489. break
  490. }
  491. }
  492. return frames
  493. }
  494. async function readControlBaseline(remote: FixtureSessionRemote): Promise<Extract<FixtureControlFrame, { type: 'baseline' }>> {
  495. const abort = new AbortController()
  496. for await (const frame of remote.control(abort.signal)) {
  497. if (frame.type !== 'baseline') continue
  498. abort.abort()
  499. return frame
  500. }
  501. throw new Error('fixture control baseline missing')
  502. }
  503. function isRemoteEventRequest(frame: FixtureRemoteEventFrame): frame is FixtureRemoteEventRequestFrame {
  504. return frame.type === 'waterfall'
  505. }
  506. function isRemoteEventCancellation(frame: FixtureRemoteEventFrame): frame is FixtureRemoteEventCancellationFrame {
  507. return frame.type === 'cancel'
  508. }
  509. async function readResidentRemoteEvents(
  510. api: FixtureTestApi,
  511. count: number,
  512. ): Promise<FixtureRemoteEventRequestFrame[]> {
  513. const abort = new AbortController()
  514. const frames = await collectValues(
  515. api.remoteEvents(abort.signal),
  516. abort,
  517. seen => seen.filter(isRemoteEventRequest).length >= count,
  518. )
  519. return frames.filter(isRemoteEventRequest)
  520. }
  521. async function nextRemoteEvent(
  522. iterator: AsyncIterator<FixtureRemoteEventFrame>,
  523. predicate: (frame: FixtureRemoteEventFrame) => boolean,
  524. ): Promise<FixtureRemoteEventFrame> {
  525. for (;;) {
  526. const item = await iterator.next()
  527. if (item.done) throw new Error('fixture Remote Event stream ended before the expected frame')
  528. if (predicate(item.value)) return item.value
  529. }
  530. }
  531. async function readOpeningCursor(remote: FixtureSessionRemote, sessionId: SessionId): Promise<number> {
  532. const abort = new AbortController()
  533. for await (const frame of remote.follow(sessionId, abort.signal)) {
  534. if (frame.type !== 'snapshot') continue
  535. abort.abort()
  536. return frame.cursor
  537. }
  538. throw new Error('fixture follow opening cursor missing')
  539. }
  540. async function readWorkspaceBaseline(
  541. remote: FixtureWorkspaceRemote,
  542. ): Promise<Extract<FixtureWorkspaceFrame, { type: 'baseline' }>['value']> {
  543. const abort = new AbortController()
  544. for await (const frame of remote.follow(abort.signal)) {
  545. if (frame.type !== 'baseline') continue
  546. abort.abort()
  547. return frame.value
  548. }
  549. throw new Error('fixture Workspace baseline missing')
  550. }
  551. describe('createFixtureApi', () => {
  552. it('keeps its Assistant frame vocabulary identical to the controller wire', () => {
  553. expectTypeOf<FixtureAssistantStreamFrame>().toEqualTypeOf<SessionAssistantStreamFrame>()
  554. })
  555. it('serves the session list sorted by updatedAt desc and echoes rpcIds on every unary', async () => {
  556. const api = createFixtureApi()
  557. const request = req({})
  558. const response = await api.sessions.list(request)
  559. expect(response.rpcId).toBe(request.rpcId)
  560. if (!response.result.ok) throw new Error('list failed')
  561. expect(response.result.value.items.map(s => s.sessionId)).toEqual(['fx-alpha', 'fx-beta', 'fx-gamma'])
  562. expect(response.result.value.items[1]?.parentSessionId).toBe('fx-alpha') // lineage material
  563. })
  564. it('returns an empty Assistant baseline when Session follow opts in', async () => {
  565. const api = createFixtureApi()
  566. const abort = new AbortController()
  567. const iterator = api.sessionRemote.follow(sid('fx-alpha'), abort.signal)[Symbol.asyncIterator]()
  568. try {
  569. const opening = await iterator.next()
  570. if (opening.done || opening.value.type !== 'snapshot') throw new Error('follow opening snapshot missing')
  571. expect(opening.value.assistantStream).toEqual({ revision: 0 })
  572. } finally {
  573. abort.abort()
  574. await iterator.return?.()
  575. }
  576. })
  577. it('searches current message text with literal unicode61-style token phrases', async () => {
  578. const api = createFixtureApi()
  579. const signal = new AbortController().signal
  580. const phrase = await api.sessions.search(req({ query: 'FIXTURE 历史消息' }), signal)
  581. expect(phrase.result).toMatchObject({
  582. ok: true,
  583. value: {
  584. items: [{ sessionId: 'fx-alpha' }],
  585. hasMore: false,
  586. },
  587. })
  588. if (!phrase.result.ok) throw new Error('search failed')
  589. expect(phrase.result.value.items[0]?.snippet).toContain('fixture 历史消息')
  590. timing().appendUser(
  591. 'fx-alpha',
  592. `${'leading context '.repeat(20)}late café token${' trailing context'.repeat(20)}`,
  593. )
  594. const late = await api.sessions.search(req({ query: 'LATE CAFE TOKEN' }), signal)
  595. if (!late.result.ok) throw new Error('late search failed')
  596. const lateSnippet = late.result.value.items[0]?.snippet ?? ''
  597. expect(lateSnippet).toContain('late café token')
  598. expect(lateSnippet.startsWith('…')).toBe(true)
  599. expect(lateSnippet.endsWith('…')).toBe(true)
  600. expect(Array.from(lateSnippet).length).toBeLessThanOrEqual(120)
  601. timing().appendUser('fx-alpha', 'Greek final sigma: ος')
  602. const finalSigma = await api.sessions.search(req({ query: 'ΟΣ' }), signal)
  603. if (!finalSigma.result.ok) throw new Error('final sigma search failed')
  604. expect(finalSigma.result.value.items[0]?.snippet).toContain('ος')
  605. const substring = await api.sessions.search(req({ query: 'ixtur' }), signal)
  606. expect(substring.result).toEqual({
  607. ok: true,
  608. value: { items: [], hasMore: false },
  609. })
  610. const punctuationOnly = await api.sessions.search(req({ query: '*' }), signal)
  611. expect(punctuationOnly.result).toEqual({
  612. ok: true,
  613. value: { items: [], hasMore: false },
  614. })
  615. const reasoningOnly = await api.sessions.search(req({ query: '思考过程' }), signal)
  616. expect(reasoningOnly.result).toEqual({
  617. ok: true,
  618. value: { items: [], hasMore: false },
  619. })
  620. const aborted = new AbortController()
  621. aborted.abort()
  622. await expect(api.sessions.search(req({ query: 'fixture' }), aborted.signal))
  623. .resolves.toMatchObject({ result: { ok: false, error: { code: 'gateway/cancelled' } } })
  624. })
  625. it('pages history backwards on message-boundary cuts with seq-contiguous stitching', async () => {
  626. const api = createFixtureApi()
  627. const tail = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 10 }))
  628. if (!tail.result.ok) throw new Error('history failed')
  629. const tailPage = tail.result.value
  630. expect(tailPage.hasMore).toBe(true)
  631. const tailEvents = historyEvents(tailPage.records)
  632. expect(tailEvents[0]?.type).toBe('turn/start') // cut lands on a turn boundary
  633. const boundary = tailEvents[0]?.seq ?? 0
  634. expect(boundary).toBeGreaterThan(0)
  635. const older = await api.sessions.history(req({ sessionId: sid('fx-alpha'), beforeSeq: boundary, maxMessages: 10 }))
  636. if (!older.result.ok) throw new Error('older failed')
  637. const olderTail = historyEvents(older.result.value.records).at(-1)
  638. expect((olderTail?.seq ?? -1) + 1).toBe(boundary) // pages stitch with no hole/overlap
  639. // Out-of-range beforeSeq clamps instead of exploding.
  640. const clamped = await api.sessions.history(req({ sessionId: sid('fx-alpha'), beforeSeq: -5, maxMessages: 10 }))
  641. if (!clamped.result.ok) throw new Error('clamped failed')
  642. expect(clamped.result.value.records).toEqual([])
  643. // Unknown session: empty page, not an error (history of a bare id).
  644. const empty = await api.sessions.history(req({ sessionId: sid('no-such'), maxMessages: 10 }))
  645. if (!empty.result.ok) throw new Error('empty failed')
  646. expect(empty.result.value).toEqual({ records: [], hasMore: false })
  647. })
  648. it('serves raw history entries with replayable tool-result metadata', async () => {
  649. const api = createFixtureApi()
  650. const response = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 200 }))
  651. if (!response.result.ok) throw new Error('history failed')
  652. const records = response.result.value.records
  653. const results = historyEvents(records)
  654. .filter(event => event.type === 'tool/result')
  655. expect(results.find(event => event.data.turn === 64)).toMatchObject({
  656. data: {
  657. meta: {
  658. diffs: [
  659. { path: 'src/config.ts', oldText: 'const timeout = 30', newText: 'const timeout = 60' },
  660. { path: 'src/config.ts', oldText: 'retries: 1', newText: 'retries: 3' },
  661. ],
  662. },
  663. },
  664. })
  665. expect(results.find(event => event.data.turn === 67)).toMatchObject({
  666. data: { meta: { shape: 'matches', truncated: true, total: 42 } },
  667. })
  668. expect(results.find(event => event.data.turn === 69)).toMatchObject({
  669. data: { meta: { path: 'packages/client/ui-primitives/src/ReadBlock.tsx', offset: 41, totalLines: 180 } },
  670. })
  671. const webSearch = results.find(event => event.data.turn === 70)
  672. expect(webSearch).toHaveProperty('data.meta.truncated', true)
  673. expect(webSearch).toHaveProperty('data.meta.sources', expect.arrayContaining([
  674. expect.objectContaining({ url: 'https://github.com/deepseek-ai/deepseek-harness' }),
  675. ]))
  676. expect(results.find(event => event.data.turn === 71)).toMatchObject({
  677. data: { meta: { url: 'https://www.deepseek.com/blog/harness-architecture', statusCode: 200 } },
  678. })
  679. const terminal = results.find(event => event.data.turn === 66)
  680. expect(terminal).toHaveProperty('data.message.content.0.content.0.type', 'text')
  681. expect(terminal).toHaveProperty(
  682. 'data.message.content.0.content.0.text',
  683. expect.stringContaining('\n[exit code: 1]'),
  684. )
  685. })
  686. it('serves grouped models and keeps a selection for later history and fixture requests', async () => {
  687. const api = createFixtureApi()
  688. const sessionId = sid('fx-alpha')
  689. const catalog = await api.sessionRemote.modelCatalog()
  690. if (!catalog.ok) throw new Error('models failed')
  691. expect(catalog.value.groups.map(group => group.name)).toEqual(['DeepSeek', 'OpenAI'])
  692. expect(catalog.value.groups[0]?.models.map(model => model.id))
  693. .toEqual(['deepseek-v4-flash', 'deepseek-v4-pro'])
  694. const selected = await api.sessions.selectModel(req({
  695. sessionId,
  696. provider: 'openai',
  697. model: 'gpt-5',
  698. }))
  699. if (!selected.result.ok) throw new Error('selection failed')
  700. expect(selected.result.value.selected).toEqual({ provider: 'openai', model: 'gpt-5' })
  701. const history = await api.sessions.history(req({ sessionId }))
  702. if (!history.result.ok) throw new Error('history failed')
  703. const prompt = await api.sessions.prompt(req({
  704. sessionId,
  705. mode: 'queue',
  706. content: [{ type: 'text', text: 'report model' }],
  707. }))
  708. expect(prompt.result.ok).toBe(true)
  709. await new Promise(resolve => setTimeout(resolve, 600))
  710. const after = await api.sessions.history(req({ sessionId }))
  711. if (!after.result.ok) throw new Error('history failed')
  712. expect(JSON.stringify(after.result.value.records)).toContain('openai/gpt-5')
  713. })
  714. it('serves configured DeepSeek readiness and keeps credential values write-only', async () => {
  715. const api = createFixtureApi()
  716. const settings = await api.settingsRemote.describe()
  717. if (!settings.ok) throw new Error('settings describe failed')
  718. expect((settings.value as { namespaces: unknown[] }).namespaces).toMatchObject([{
  719. ns: 'llm-deepseek',
  720. value: { apiKeyEnv: 'DEEPSEEK_API_KEY' },
  721. secrets: [{ path: ['apiKey'], set: false }],
  722. }])
  723. for (const result of [
  724. await api.settingsRemote.update('llm-deepseek', {}, undefined),
  725. await api.settingsRemote.replace('llm-deepseek', {}, undefined),
  726. ]) {
  727. expect(result).toMatchObject({
  728. ok: false,
  729. error: { code: 'settings/rejected', message: 'fixture: the minimal readiness settings descriptor is read-only' },
  730. })
  731. }
  732. const describe = async (refs: readonly string[]): Promise<Record<string, unknown>> => {
  733. const result = await api.credentialRemote.describe(refs)
  734. if (!result.ok) throw new Error('credential describe failed')
  735. return result.value as Record<string, unknown>
  736. }
  737. expect(await describe(['DEEPSEEK_API_KEY', 'TEST_API_KEY'])).toEqual({
  738. DEEPSEEK_API_KEY: { configured: true, source: 'file', writable: true },
  739. TEST_API_KEY: { configured: false, writable: true },
  740. })
  741. await api.credentialRemote.set('TEST_API_KEY', 'write-only-fixture-secret')
  742. expect((await describe(['TEST_API_KEY'])).TEST_API_KEY).toEqual({
  743. configured: true,
  744. source: 'file',
  745. writable: true,
  746. })
  747. await api.credentialRemote.unset('TEST_API_KEY')
  748. expect((await describe(['TEST_API_KEY'])).TEST_API_KEY).toEqual({ configured: false, writable: true })
  749. })
  750. it('emits the todo/write snapshot at the real tool boundary: between tool/call and tool/result, timestamps monotonic', async () => {
  751. const api = createFixtureApi()
  752. const tail = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 10 }))
  753. if (!tail.result.ok) throw new Error('history failed')
  754. const events = historyEvents(tail.result.value.records)
  755. const todoAt = events.findIndex(e => e.type === 'todo/write')
  756. expect(todoAt).toBeGreaterThan(0)
  757. // Production ordering (the tool appends mid-execution): call → snapshot → result.
  758. expect(events[todoAt - 1]?.type).toBe('tool/call')
  759. expect(events[todoAt + 1]?.type).toBe('tool/result')
  760. const times = events.slice(todoAt - 1, todoAt + 2).map(e => e.time)
  761. expect(times[0]).toBeLessThanOrEqual(times[1] ?? 0)
  762. expect(times[1]).toBeLessThanOrEqual(times[2] ?? 0)
  763. // The sample is a parallel plan: this fixture chooses the parallel policy,
  764. // so the surfaces fed from here face more than one active item.
  765. const snapshot = events[todoAt] as { data: { todos: { status: string }[] } }
  766. expect(snapshot.data.todos.filter(t => t.status === 'in_progress')).toHaveLength(2)
  767. })
  768. it('create adds a session and announces it through the Host Remote event stream', async () => {
  769. const api = createFixtureApi()
  770. const abort = new AbortController()
  771. const seen: FixtureRemoteEventNotificationFrame[] = []
  772. const consuming = (async () => {
  773. for await (const frame of api.remoteEvents(abort.signal)) {
  774. if (frame.type !== 'emit' || frame.event !== 'api-session/added') continue
  775. seen.push(frame)
  776. abort.abort()
  777. break
  778. }
  779. })()
  780. await new Promise(resolve => setTimeout(resolve, 10)) // let the stream register
  781. const created = await api.sessions.create(req({}))
  782. if (!created.result.ok) throw new Error('create failed')
  783. await consuming
  784. if (!created.result.ok) throw new Error('create failed')
  785. const createdId = created.result.value.sessionId
  786. expect(seen).toHaveLength(1)
  787. const added = seen[0]
  788. expect(added).toMatchObject({
  789. event: 'api-session/added',
  790. args: [{ sessionId: createdId, blank: true, cwd: '/tmp/fixture' }],
  791. })
  792. const list = await api.sessions.list(req({}))
  793. if (!list.result.ok) throw new Error('list failed')
  794. expect(list.result.value.items.some(s => s.sessionId === createdId)).toBe(true)
  795. })
  796. it('prompt replays a full streamed turn and cancel mid-replay freezes with (已中断)', async () => {
  797. const api = createFixtureApi()
  798. const created = await api.sessions.create(req({}))
  799. if (!created.result.ok) throw new Error('create failed')
  800. const id = created.result.value.sessionId
  801. const followAbort = new AbortController()
  802. const controlAbort = new AbortController()
  803. const controlFrames: FixtureControlFrame[] = []
  804. const followPromise = collectValues(
  805. api.sessionRemote.follow(id, followAbort.signal),
  806. followAbort,
  807. frames => frames.some(frame => frame.type === 'event' && frame.event.type === 'turn/end'),
  808. )
  809. const controlPromise = (async () => {
  810. for await (const frame of api.sessionRemote.control(controlAbort.signal)) controlFrames.push(frame)
  811. })()
  812. await new Promise(resolve => setTimeout(resolve, 10))
  813. // Unknown session → session/not-found with the id echoed in details.
  814. const missing = await api.sessions.prompt(req({ sessionId: sid('ghost'), mode: 'queue' as const, content: [{ type: 'text' as const, text: 'x' }] }))
  815. expect(missing.result).toMatchObject({ ok: false, error: { code: 'session/not-found', details: { sessionId: 'ghost' } } })
  816. // Real prompt: replay starts (running flips true), cancel freezes it.
  817. const accepted = await api.sessions.prompt(req({ sessionId: id, mode: 'queue' as const, content: [{ type: 'text' as const, text: 'render markdown' }] }))
  818. expect(accepted.result).toMatchObject({ ok: true, value: { accepted: true } })
  819. await new Promise(resolve => setTimeout(resolve, 120)) // a couple of typewriter ticks
  820. await api.sessions.cancel(req({ sessionId: id }))
  821. const frames = await followPromise
  822. const types = frames.flatMap(frame => frame.type === 'event' ? [frame.event.type] : [])
  823. const assistantFrames = frames.flatMap(frame => frame.type === 'assistant-stream' ? [frame.frame] : [])
  824. expect(types).toContain('turn/start')
  825. expect(types).toContain('user/message')
  826. expect(types).toContain('assistant/message')
  827. expect(assistantFrames.some(frame => frame.type === 'chunk')).toBe(true)
  828. const committed = assistantFrames.find(frame => frame.type === 'end'
  829. && frame.outcome.kind === 'committed'
  830. && frame.outcome.eventType === 'assistant/message')
  831. expect(committed).toMatchObject({
  832. type: 'end',
  833. index: assistantFrames.filter(frame => frame.type === 'chunk'
  834. && frame.attemptId === committed?.attemptId).length,
  835. })
  836. expect(types.at(-1)).toBe('turn/end')
  837. // Capacity is durable log state, not a transient frame: the prompt path
  838. // records request/context and the projection carries it to the client.
  839. expect(types).toContain('request/context')
  840. await vi.waitFor(() => {
  841. expect(controlFrames.some(frame =>
  842. frame.type === 'projection'
  843. && frame.key === 'contextBreakdown'
  844. && (frame.value as { messageTokens?: number }).messageTokens! > 0)).toBe(true)
  845. })
  846. expect(controlFrames.some(frame =>
  847. frame.type === 'projection'
  848. && frame.key === 'tokenUsage'
  849. && (frame.value as { outputTokens?: number }).outputTokens === 8)).toBe(true)
  850. expect(controlFrames.some(frame =>
  851. frame.type === 'projection'
  852. && frame.key === 'contextPressure'
  853. && (frame.value as { contextWindow?: number }).contextWindow === 128_000)).toBe(true)
  854. const finalize = frames.find(frame => frame.type === 'event' && frame.event.type === 'assistant/message')
  855. if (finalize?.type !== 'event') throw new Error('assistant final event missing')
  856. if (finalize.event.type !== 'assistant/message') throw new Error('assistant final event has wrong type')
  857. expect(finalize.event.data.interrupted).toBe(true)
  858. expect(finalize.event.data.stream.length).toBeGreaterThan(0)
  859. controlAbort.abort()
  860. await controlPromise
  861. // Idle cancel: no replay in flight, must not explode; running flips false.
  862. const idleCancel = await api.sessions.cancel(req({ sessionId: id }))
  863. expect(idleCancel.result).toMatchObject({ ok: true })
  864. })
  865. it('steer during a replay lands a user/message inside the current turn and the replay continues', async () => {
  866. const api = createFixtureApi()
  867. const created = await api.sessions.create(req({}))
  868. if (!created.result.ok) throw new Error('create failed')
  869. const id = created.result.value.sessionId
  870. const abort = new AbortController()
  871. const framesPromise = collectValues(api.sessionRemote.follow(id, abort.signal), abort,
  872. frames => frames.some(frame => frame.type === 'event' && frame.event.type === 'turn/end'))
  873. await new Promise(resolve => setTimeout(resolve, 10))
  874. await api.sessions.prompt(req({ sessionId: id, mode: 'queue' as const, content: [{ type: 'text' as const, text: '短' }] }))
  875. await api.sessions.prompt(req({ sessionId: id, mode: 'steer' as const, content: [{ type: 'text' as const, text: '插话' }] }))
  876. const frames = await framesPromise
  877. const types = frames.flatMap(frame => frame.type === 'event' ? [frame.event.type] : [])
  878. expect(JSON.stringify(frames)).toContain('插话')
  879. expect(types.at(-1)).toBe('turn/end') // steer did not restart the turn
  880. })
  881. it('control replays projections while resident Remote Events retain ids across reconnects', async () => {
  882. const api = createFixtureApi()
  883. const first = await readControlBaseline(api.sessionRemote)
  884. const second = await readControlBaseline(api.sessionRemote)
  885. expect(first.value.approvals).toEqual([])
  886. expect(first.value.questions).toEqual([])
  887. const alpha = first.value.projections['fx-alpha']
  888. expect(alpha?.asOfSeq).toBeGreaterThan(0)
  889. expect(alpha?.values).toMatchObject({
  890. title: 'Fixture 历史会话',
  891. plan: { active: false, pending: false },
  892. goal: null,
  893. imageLimits: { maxImagesPerMessage: 20, maxImageBytes: 5 * 1024 * 1024 },
  894. })
  895. const breakdown = alpha?.values['contextBreakdown'] as { systemTokens: number; messageTokens: number }
  896. expect(breakdown.messageTokens).toBeGreaterThan(0)
  897. // The seeded system/message at surface node 0 prices the system figure.
  898. expect(breakdown.systemTokens).toBeGreaterThan(0)
  899. expect((alpha?.values['sessionStats'] as { steps: number }).steps).toBeGreaterThan(0)
  900. expect(second.value.projections['fx-alpha']).toEqual(alpha)
  901. const firstEvents = await readResidentRemoteEvents(api, 2)
  902. const secondEvents = await readResidentRemoteEvents(api, 2)
  903. const firstApproval = firstEvents.find(frame => frame.event === 'approval/request')
  904. const firstQuestion = firstEvents.find(frame => frame.event === 'user-questions/request')
  905. const secondApproval = secondEvents.find(frame => frame.event === 'approval/request')
  906. const secondQuestion = secondEvents.find(frame => frame.event === 'user-questions/request')
  907. expect(firstApproval).toMatchObject({
  908. type: 'waterfall',
  909. request: { toolName: 'dangerous_tool' },
  910. agentId: 'fx-alpha',
  911. })
  912. expect(firstQuestion).toMatchObject({
  913. type: 'waterfall',
  914. agentId: 'fx-alpha',
  915. })
  916. expect(Array.isArray(firstQuestion?.request.questions)).toBe(true)
  917. expect(secondApproval?.eventId).toBe(firstApproval?.eventId)
  918. expect(secondQuestion?.eventId).toBe(firstQuestion?.eventId)
  919. expect(await readOpeningCursor(api.sessionRemote, sid('fx-alpha'))).toBeGreaterThan(0)
  920. })
  921. it('steer with no replay in flight falls through to a fresh queued turn; non-text blocks stringify empty', async () => {
  922. const api = createFixtureApi()
  923. const created = await api.sessions.create(req({}))
  924. if (!created.result.ok) throw new Error('create failed')
  925. const abort = new AbortController()
  926. const framesPromise = collectValues(
  927. api.sessionRemote.follow(created.result.value.sessionId, abort.signal),
  928. abort,
  929. frames => frames.some(frame => frame.type === 'event' && frame.event.type === 'turn/end'),
  930. )
  931. await new Promise(resolve => setTimeout(resolve, 10))
  932. // steer while idle + a non-text content block (covers the '' arm of the text join).
  933. await api.sessions.prompt(req({
  934. sessionId: created.result.value.sessionId, mode: 'steer' as const,
  935. content: [{ type: 'text' as const, text: '短' }, { type: 'image', data: 'x' } as never],
  936. }))
  937. const frames = await framesPromise
  938. const types = frames.flatMap(frame => frame.type === 'event' ? [frame.event.type] : [])
  939. expect(types[0]).toBe('turn/start') // idle steer degraded to a queued turn, not an in-turn insert
  940. })
  941. it('gamma interval flip emits a Remote status event and its empty follow source opens at -1', async () => {
  942. vi.useFakeTimers()
  943. try {
  944. const api = createFixtureApi()
  945. const abort = new AbortController()
  946. const hostSeen: FixtureRemoteEventFrame[] = []
  947. const consuming = (async () => {
  948. for await (const frame of api.remoteEvents(abort.signal)) hostSeen.push(frame)
  949. })()
  950. await vi.advanceTimersByTimeAsync(5001) // interval fires: fx-gamma flips running=true (no log exists)
  951. expect(hostSeen).toContainEqual({
  952. type: 'emit',
  953. event: 'api-session/status',
  954. args: [sid('fx-gamma'), true],
  955. })
  956. expect(await readOpeningCursor(api.sessionRemote, sid('fx-gamma'))).toBe(-1)
  957. abort.abort()
  958. await vi.advanceTimersByTimeAsync(10)
  959. await consuming
  960. } finally {
  961. vi.useRealTimers()
  962. }
  963. })
  964. it('answers a resident question through its Remote Event id and stops replaying it', async () => {
  965. const api = createFixtureApi()
  966. const abort = new AbortController()
  967. const stream = api.remoteEvents(abort.signal)
  968. const iterator = stream[Symbol.asyncIterator]()
  969. const question = await nextRemoteEvent(
  970. iterator,
  971. frame => isRemoteEventRequest(frame) && frame.event === 'user-questions/request',
  972. )
  973. if (!isRemoteEventRequest(question)) throw new Error('fixture question Remote Event missing')
  974. const clientId = await stream.clientId
  975. await expect(api.answerRemoteEvent({
  976. clientId,
  977. eventId: 'unrelated',
  978. outcome: { kind: 'result', value: {} },
  979. })).resolves.toEqual({ ok: true, value: undefined })
  980. await expect(api.answerRemoteEvent({
  981. clientId,
  982. eventId: question.eventId,
  983. outcome: { kind: 'result', value: { answers: {} } },
  984. })).resolves.toEqual({ ok: true, value: undefined })
  985. const cancelled = await nextRemoteEvent(
  986. iterator,
  987. frame => isRemoteEventCancellation(frame) && frame.eventId === question.eventId,
  988. )
  989. expect(cancelled).toEqual({ type: 'cancel', eventId: question.eventId })
  990. abort.abort()
  991. await iterator.return?.()
  992. await expect(api.answerRemoteEvent({
  993. clientId,
  994. eventId: question.eventId,
  995. outcome: { kind: 'result', value: { answers: {} } },
  996. })).resolves.toMatchObject({ ok: false, error: { code: 'gateway/invocation-unavailable' } })
  997. const remaining = await readResidentRemoteEvents(api, 1)
  998. expect(remaining.map(frame => frame.event)).toEqual(['approval/request'])
  999. const cancelledApi = createFixtureApi()
  1000. const cancelAbort = new AbortController()
  1001. const cancelStream = cancelledApi.remoteEvents(cancelAbort.signal)
  1002. const cancelIterator = cancelStream[Symbol.asyncIterator]()
  1003. const cancelQuestion = await nextRemoteEvent(
  1004. cancelIterator,
  1005. frame => isRemoteEventRequest(frame) && frame.event === 'user-questions/request',
  1006. )
  1007. if (!isRemoteEventRequest(cancelQuestion)) throw new Error('fixture cancellation question missing')
  1008. await expect(cancelledApi.answerRemoteEvent({
  1009. clientId: await cancelStream.clientId,
  1010. eventId: cancelQuestion.eventId,
  1011. outcome: {
  1012. kind: 'rejected',
  1013. error: { name: 'UserQuestionError', message: 'skip', code: 'ASK_CANCELLED' },
  1014. },
  1015. })).resolves.toEqual({ ok: true, value: undefined })
  1016. cancelAbort.abort()
  1017. await cancelIterator.return?.()
  1018. const afterCancellation = await readResidentRemoteEvents(cancelledApi, 1)
  1019. expect(afterCancellation.map(frame => frame.event)).toEqual(['approval/request'])
  1020. })
  1021. it('answers a resident approval and broadcasts cancellation to its active delivery', async () => {
  1022. const api = createFixtureApi()
  1023. const abort = new AbortController()
  1024. const stream = api.remoteEvents(abort.signal)
  1025. const iterator = stream[Symbol.asyncIterator]()
  1026. const approval = await nextRemoteEvent(
  1027. iterator,
  1028. frame => isRemoteEventRequest(frame) && frame.event === 'approval/request',
  1029. )
  1030. if (!isRemoteEventRequest(approval)) throw new Error('fixture approval Remote Event missing')
  1031. await expect(api.answerRemoteEvent({
  1032. clientId: await stream.clientId,
  1033. eventId: approval.eventId,
  1034. outcome: { kind: 'result', value: 'allowed-once' },
  1035. })).resolves.toEqual({ ok: true, value: undefined })
  1036. const cancelled = await nextRemoteEvent(
  1037. iterator,
  1038. frame => isRemoteEventCancellation(frame) && frame.eventId === approval.eventId,
  1039. )
  1040. expect(cancelled).toEqual({ type: 'cancel', eventId: approval.eventId })
  1041. abort.abort()
  1042. await iterator.return?.()
  1043. await expect(api.answerRemoteEvent({
  1044. clientId: await stream.clientId,
  1045. eventId: approval.eventId,
  1046. outcome: { kind: 'next' },
  1047. })).resolves.toMatchObject({ ok: false, error: { code: 'gateway/invocation-unavailable' } })
  1048. const remaining = await readResidentRemoteEvents(api, 1)
  1049. expect(remaining.map(frame => frame.event)).toEqual(['user-questions/request'])
  1050. })
  1051. it('createDirectory under the root mints /name whose listing and crumbs share the identity', async () => {
  1052. const api = createFixtureApi()
  1053. const created = await api.directoryPickerRemote.createDirectory('/', 'srv')
  1054. if (!created.ok) throw new Error('create failed')
  1055. expect(created.value).toBe('/srv')
  1056. const listed = await api.directoryPickerRemote.list('/srv')
  1057. if (!listed.ok) throw new Error('list failed')
  1058. expect(listed.value.crumbs).toEqual([
  1059. { name: '/', path: '/', hidden: false },
  1060. { name: 'srv', path: '/srv', hidden: false },
  1061. ])
  1062. const root = await api.directoryPickerRemote.list('/')
  1063. if (!root.ok) throw new Error('root list failed')
  1064. expect(root.value.entries).toContainEqual({ name: 'srv', path: '/srv', hidden: false })
  1065. })
  1066. it('workspace/follow serves the resident baseline and create reuses on path collision', async () => {
  1067. const api = createFixtureApi()
  1068. const baseline = await readWorkspaceBaseline(api.workspaceRemote)
  1069. expect(baseline.items).toEqual([
  1070. expect.objectContaining({
  1071. workspaceId: 'fx-ws-fixture', path: '/tmp/fixture', title: 'fixture',
  1072. sessionIds: ['fx-alpha', 'fx-beta', 'fx-gamma'],
  1073. }),
  1074. expect.objectContaining({
  1075. workspaceId: 'fx-ws-home', path: '/home/fixture/Documents/project', title: 'project',
  1076. sessionIds: [],
  1077. }),
  1078. ])
  1079. // path collision → the existing entity comes back, created:false, no frame.
  1080. const reused = await api.workspace.create(req({ path: '/tmp/fixture' }))
  1081. if (!reused.result.ok) throw new Error('reuse failed')
  1082. expect(reused.result.value).toMatchObject({ created: false, workspace: { workspaceId: 'fx-ws-fixture' } })
  1083. })
  1084. it('workspace.create on a fresh path mints a new entity and pushes an upsert', async () => {
  1085. const api = createFixtureApi()
  1086. const abort = new AbortController()
  1087. const consuming = collectValues(
  1088. api.workspaceRemote.follow(abort.signal),
  1089. abort,
  1090. frames => frames.some(frame => frame.type === 'upsert'
  1091. && frame.workspace.path === '/tmp/fixture-workspaces/nova'),
  1092. )
  1093. await new Promise(resolve => setTimeout(resolve, 10))
  1094. const created = await api.workspace.create(req({ path: '/tmp/fixture-workspaces/nova' }))
  1095. if (!created.result.ok) throw new Error('create failed')
  1096. expect(created.result.value.created).toBe(true)
  1097. expect(created.result.value.workspace).toMatchObject({
  1098. path: '/tmp/fixture-workspaces/nova', title: 'nova', sessionIds: [],
  1099. })
  1100. const frames = await consuming
  1101. expect(frames.at(-1)).toEqual({ type: 'upsert', workspace: created.result.value.workspace })
  1102. // A basename-less path serves as its own title.
  1103. const rootPath = await api.workspace.create(req({ path: '/' }))
  1104. if (!rootPath.result.ok) throw new Error('rootPath failed')
  1105. expect(rootPath.result.value.workspace.title).toBe('/')
  1106. })
  1107. it('workspace.rename covers not-found, conflict, no-op, and the changed frame', async () => {
  1108. const api = createFixtureApi()
  1109. const abort = new AbortController()
  1110. const consuming = collectValues(
  1111. api.workspaceRemote.follow(abort.signal),
  1112. abort,
  1113. frames => frames.filter(frame => frame.type === 'upsert').length >= 2,
  1114. )
  1115. await new Promise(resolve => setTimeout(resolve, 10))
  1116. const wsid = 'fx-ws-fixture' as WorkspaceId
  1117. const missing = await api.workspace.rename(req({ workspaceId: 'fx-ws-void' as WorkspaceId, title: 'x' }))
  1118. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace/not-found', details: { workspaceId: 'fx-ws-void' } } })
  1119. await api.workspace.create(req({ path: '/tmp/fixture-workspaces/occupied' }))
  1120. const conflict = await api.workspace.rename(req({ workspaceId: wsid, title: ' occupied ' }))
  1121. expect(conflict.result).toMatchObject({ ok: false, error: { code: 'workspace/name-conflict', details: { name: 'occupied' } } })
  1122. const noop = await api.workspace.rename(req({ workspaceId: wsid, title: ' fixture ' }))
  1123. if (!noop.result.ok) throw new Error('no-op rename failed')
  1124. expect(noop.result.value.workspace.title).toBe('fixture')
  1125. const renamed = await api.workspace.rename(req({ workspaceId: wsid, title: 'renamed' }))
  1126. if (!renamed.result.ok) throw new Error('rename failed')
  1127. expect(renamed.result.value.workspace.title).toBe('renamed')
  1128. const frames = await consuming
  1129. // Only the create and the effective rename emit frames; the no-op stays silent.
  1130. const upserts = frames.filter(frame => frame.type === 'upsert')
  1131. expect(upserts).toHaveLength(2)
  1132. expect(upserts[1]).toMatchObject({ workspace: { workspaceId: wsid, title: 'renamed' } })
  1133. })
  1134. it('session.rename covers not-found, blank title, and the accepted append + title frame', async () => {
  1135. const api = createFixtureApi()
  1136. const followAbort = new AbortController()
  1137. const controlAbort = new AbortController()
  1138. const followPromise = collectValues(
  1139. api.sessionRemote.follow(sid('fx-alpha'), followAbort.signal),
  1140. followAbort,
  1141. frames => frames.some(frame => frame.type === 'event'
  1142. && (frame.event as { type: string }).type === 'session/title'),
  1143. )
  1144. const controlPromise = collectValues(
  1145. api.sessionRemote.control(controlAbort.signal),
  1146. controlAbort,
  1147. frames => frames.some(frame => frame.type === 'projection' && frame.key === 'title' && frame.value === '重命名'),
  1148. )
  1149. await new Promise(resolve => setTimeout(resolve, 10))
  1150. const missing = await api.sessions.rename(req({ sessionId: sid('fx-void'), title: 'x' }))
  1151. expect(missing.result).toMatchObject({ ok: false, error: { code: 'session/not-found', details: { sessionId: 'fx-void' } } })
  1152. const blank = await api.sessions.rename(req({ sessionId: sid('fx-alpha'), title: ' ' }))
  1153. expect(blank.result).toMatchObject({ ok: false, error: { code: 'session/title-invalid', details: { sessionId: 'fx-alpha' } } })
  1154. const renamed = await api.sessions.rename(req({ sessionId: sid('fx-alpha'), title: ' 重命名 ' }))
  1155. if (!renamed.result.ok) throw new Error('rename failed')
  1156. expect(renamed.result.value.title).toBe('重命名')
  1157. const acceptedSeq = renamed.result.value.seq
  1158. // The response seq addresses the appended title event (the client plane
  1159. // has no session/title in its event union — titles ride the projection —
  1160. // so the event is located by seq and its payload checked structurally).
  1161. const history = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 100 }))
  1162. if (!history.result.ok) throw new Error('history failed')
  1163. const appended = historyEvents(history.result.value.records).find(event => event.seq === acceptedSeq)
  1164. expect(appended).toMatchObject({
  1165. type: 'session/title',
  1166. data: { title: '重命名', messageSeqs: [], source: { kind: 'user' } },
  1167. })
  1168. const followed = await followPromise
  1169. expect(followed.some(frame => frame.type === 'event'
  1170. && frame.event.seq === acceptedSeq
  1171. && (frame.event as { readonly type: string }).type === 'session/title')).toBe(true)
  1172. const frames = await controlPromise
  1173. const titleFrames = frames.filter(frame =>
  1174. frame.type === 'projection'
  1175. && frame.key === 'title'
  1176. && frame.sessionId === sid('fx-alpha')
  1177. && frame.value === '重命名')
  1178. expect(titleFrames).toHaveLength(1)
  1179. expect(titleFrames[0]).toMatchObject({ seq: acceptedSeq })
  1180. })
  1181. it('workspace.insertSessionBefore moves, appends, no-ops, and rejects invalid ids', async () => {
  1182. const api = createFixtureApi()
  1183. const wsid = 'fx-ws-fixture' as WorkspaceId
  1184. const missing = await api.workspace.insertSessionBefore(req({ workspaceId: 'fx-ws-void' as WorkspaceId, sessionId: sid('fx-alpha') }))
  1185. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace/not-found' } })
  1186. const ghost = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-ghost') }))
  1187. expect(ghost.result).toMatchObject({ ok: false, error: { code: 'workspace/move-invalid', details: { sessionId: 'fx-ghost' } } })
  1188. const badAnchor = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-alpha'), beforeSessionId: sid('fx-ghost') }))
  1189. expect(badAnchor.result).toMatchObject({ ok: false, error: { code: 'workspace/move-invalid', details: { beforeSessionId: 'fx-ghost' } } })
  1190. const moved = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-gamma'), beforeSessionId: sid('fx-beta') }))
  1191. if (!moved.result.ok) throw new Error('move failed')
  1192. expect(moved.result.value.workspace.sessionIds).toEqual(['fx-alpha', 'fx-gamma', 'fx-beta'])
  1193. const appended = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-alpha') }))
  1194. if (!appended.result.ok) throw new Error('append failed')
  1195. expect(appended.result.value.workspace.sessionIds).toEqual(['fx-gamma', 'fx-beta', 'fx-alpha'])
  1196. const before = appended.result.value.workspace.updatedAt
  1197. const noop = await api.workspace.insertSessionBefore(req({ workspaceId: wsid, sessionId: sid('fx-alpha') }))
  1198. if (!noop.result.ok) throw new Error('no-op move failed')
  1199. expect(noop.result.value.workspace.sessionIds).toEqual(['fx-gamma', 'fx-beta', 'fx-alpha'])
  1200. expect(noop.result.value.workspace.updatedAt).toBe(before)
  1201. })
  1202. it('workspace.delete removes only the Workspace row and emits the removal frame', async () => {
  1203. const api = createFixtureApi()
  1204. const abort = new AbortController()
  1205. const consuming = collectValues(
  1206. api.workspaceRemote.follow(abort.signal),
  1207. abort,
  1208. frames => frames.some(frame => frame.type === 'remove'),
  1209. )
  1210. await new Promise(resolve => setTimeout(resolve, 10))
  1211. const missing = await api.workspace.delete(req({ workspaceId: 'fx-ws-void' as WorkspaceId }))
  1212. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace/not-found' } })
  1213. const deleted = await api.workspace.delete(req({ workspaceId: 'fx-ws-fixture' as WorkspaceId }))
  1214. expect(deleted.result).toEqual({ ok: true, value: { deleted: true } })
  1215. const frames = await consuming
  1216. expect(frames.at(-1)).toEqual({ type: 'remove', workspaceId: 'fx-ws-fixture' })
  1217. const baseline = await readWorkspaceBaseline(api.workspaceRemote)
  1218. expect(baseline.items.some(workspace => workspace.workspaceId === 'fx-ws-fixture')).toBe(false)
  1219. const sessions = await api.sessions.list(req({}))
  1220. if (!sessions.result.ok) throw new Error('session list failed')
  1221. expect(sessions.result.value.items.map(session => session.sessionId)).toContain('fx-alpha')
  1222. })
  1223. it('session.create({workspaceId}) lands on the account and unknown ids error', async () => {
  1224. const api = createFixtureApi()
  1225. const hostAbort = new AbortController()
  1226. const workspaceAbort = new AbortController()
  1227. const seen: FixtureRemoteEventNotificationFrame[] = []
  1228. const consuming = (async () => {
  1229. for await (const frame of api.remoteEvents(hostAbort.signal)) {
  1230. if (frame.type !== 'emit' || frame.event !== 'api-session/added') continue
  1231. seen.push(frame)
  1232. hostAbort.abort()
  1233. break
  1234. }
  1235. })()
  1236. const workspaceFrames = collectValues(
  1237. api.workspaceRemote.follow(workspaceAbort.signal),
  1238. workspaceAbort,
  1239. frames => frames.some(frame => frame.type === 'upsert'
  1240. && frame.workspace.sessionIds.length === 4),
  1241. )
  1242. await new Promise(resolve => setTimeout(resolve, 10))
  1243. const missing = await api.sessions.create(req({ workspaceId: 'fx-ws-void' as WorkspaceId }))
  1244. expect(missing.result).toMatchObject({ ok: false, error: { code: 'workspace/not-found', details: { workspaceId: 'fx-ws-void' } } })
  1245. const created = await api.sessions.create(req({ workspaceId: 'fx-ws-fixture' as WorkspaceId }))
  1246. if (!created.result.ok) throw new Error('create failed')
  1247. const id = created.result.value.sessionId
  1248. await consuming
  1249. const added = seen[0]
  1250. expect(added).toMatchObject({
  1251. event: 'api-session/added',
  1252. args: [{ sessionId: id, blank: true, cwd: '/tmp/fixture' }],
  1253. })
  1254. expect((await workspaceFrames).at(-1)).toMatchObject({
  1255. type: 'upsert',
  1256. workspace: {
  1257. workspaceId: 'fx-ws-fixture',
  1258. sessionIds: [id, 'fx-alpha', 'fx-beta', 'fx-gamma'],
  1259. },
  1260. })
  1261. })
  1262. it('supports an empty baseline, preallocated ids, independent streams, and idempotent retry', async () => {
  1263. const api = createFixtureApi({ empty: true, createFrameOrder: 'workspace-first' })
  1264. const initialSessions = await api.sessions.list(req({}))
  1265. expect(initialSessions.result).toMatchObject({ ok: true, value: { items: [] } })
  1266. expect(await readWorkspaceBaseline(api.workspaceRemote)).toEqual({
  1267. items: [],
  1268. archivedSessionIds: [],
  1269. })
  1270. const made = await api.workspace.create(req({ path: '/tmp/fixture-workspaces/nova' }))
  1271. if (!made.result.ok) throw new Error('workspace create failed')
  1272. const hostAbort = new AbortController()
  1273. const workspaceAbort = new AbortController()
  1274. const hostFrames = collectValues(
  1275. api.remoteEvents(hostAbort.signal),
  1276. hostAbort,
  1277. frames => frames.length === 1,
  1278. )
  1279. const workspaceFrames = collectValues(
  1280. api.workspaceRemote.follow(workspaceAbort.signal),
  1281. workspaceAbort,
  1282. frames => frames.some(frame => frame.type === 'upsert'
  1283. && frame.workspace.sessionIds.includes(sid('fx-preallocated'))),
  1284. )
  1285. await new Promise(resolve => setTimeout(resolve, 10))
  1286. const preallocated = sid('fx-preallocated')
  1287. const created = await api.sessions.create(req({
  1288. workspaceId: made.result.value.workspace.workspaceId,
  1289. sessionId: preallocated,
  1290. }))
  1291. expect(created.result).toEqual({ ok: true, value: { sessionId: preallocated } })
  1292. expect((await workspaceFrames).at(-1)).toMatchObject({
  1293. type: 'upsert', workspace: { sessionIds: [preallocated] },
  1294. })
  1295. const added = (await hostFrames)[0]
  1296. expect(added).toMatchObject({
  1297. event: 'api-session/added',
  1298. args: [{
  1299. sessionId: preallocated,
  1300. blank: true,
  1301. cwd: made.result.value.workspace.path,
  1302. }],
  1303. })
  1304. const retried = await api.sessions.create(req({
  1305. workspaceId: made.result.value.workspace.workspaceId,
  1306. sessionId: preallocated,
  1307. }))
  1308. expect(retried.result).toEqual({ ok: true, value: { sessionId: preallocated } })
  1309. const listed = await api.sessions.list(req({}))
  1310. if (!listed.result.ok) throw new Error('session list failed')
  1311. expect(listed.result.value.items.filter(item => item.sessionId === preallocated)).toHaveLength(1)
  1312. const conflict = await api.sessions.create(req({ sessionId: preallocated, cwd: '/elsewhere' }))
  1313. expect(conflict.result).toMatchObject({
  1314. ok: false,
  1315. error: { code: 'session/conflict', details: { sessionId: preallocated, requestedCwd: '/elsewhere' } },
  1316. })
  1317. })
  1318. it('attaches an existing ungrouped Session to a matching Workspace', async () => {
  1319. const api = createFixtureApi()
  1320. const sessionId = sid('fx-existing-ungrouped')
  1321. await expect(api.sessions.create(req({ sessionId, cwd: '/tmp/fixture' }))).resolves.toMatchObject({
  1322. result: { ok: true, value: { sessionId } },
  1323. })
  1324. await expect(api.sessions.create(req({
  1325. sessionId,
  1326. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1327. }))).resolves.toMatchObject({ result: { ok: true, value: { sessionId } } })
  1328. const workspaces = await readWorkspaceBaseline(api.workspaceRemote)
  1329. expect(workspaces.items[0]?.sessionIds).toContain(sessionId)
  1330. })
  1331. it('reports a conflict without an existing cwd detail for an unrecorded cwd', async () => {
  1332. const api = createFixtureApi()
  1333. const listed = await api.sessions.list(req({}))
  1334. if (!listed.result.ok) throw new Error('session list failed')
  1335. const existing = listed.result.value.items.find(item => item.sessionId === sid('fx-alpha'))
  1336. if (existing === undefined) throw new Error('fixture Session missing')
  1337. delete existing.cwd
  1338. const conflict = await api.sessions.create(req({ sessionId: existing.sessionId }))
  1339. expect(conflict.result).toEqual({
  1340. ok: false,
  1341. error: {
  1342. code: 'session/conflict',
  1343. message: `session ${existing.sessionId} already uses no cwd`,
  1344. details: { sessionId: existing.sessionId, requestedCwd: '/tmp/fixture' },
  1345. },
  1346. })
  1347. })
  1348. it('publishes an ungrouped Session when Workspace attachment fails', async () => {
  1349. const api = createFixtureApi({ failWorkspaceAttach: true })
  1350. const sessionId = sid('fx-partial')
  1351. const created = await api.sessions.create(req({
  1352. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1353. sessionId,
  1354. }))
  1355. expect(created.result).toMatchObject({
  1356. ok: false,
  1357. error: { code: 'session/workspace-attach-failed', details: { sessionId, workspaceId: 'fx-ws-fixture' } },
  1358. })
  1359. const listed = await api.sessions.list(req({}))
  1360. const workspaces = await readWorkspaceBaseline(api.workspaceRemote)
  1361. if (!listed.result.ok) throw new Error('list failed')
  1362. expect(listed.result.value.items.filter(item => item.sessionId === sessionId)).toHaveLength(1)
  1363. expect(workspaces.items[0]?.sessionIds).not.toContain(sessionId)
  1364. const retried = await api.sessions.create(req({
  1365. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1366. sessionId,
  1367. }))
  1368. expect(retried.result).toMatchObject({ ok: false, error: { code: 'session/workspace-attach-failed' } })
  1369. const afterRetry = await api.sessions.list(req({}))
  1370. if (!afterRetry.result.ok) throw new Error('list failed')
  1371. expect(afterRetry.result.value.items.filter(item => item.sessionId === sessionId)).toHaveLength(1)
  1372. })
  1373. it('reconciles a dropped create response and can reject a prompt before acceptance', async () => {
  1374. const sessionId = sid('fx-lost-response')
  1375. const dropped = createFixtureApi({ dropSessionCreateResponse: true })
  1376. await expect(Promise.resolve().then(() => dropped.sessions.create(req({
  1377. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1378. sessionId,
  1379. })))).rejects.toThrow(/dropped session\.create response/)
  1380. const listed = await dropped.sessions.list(req({}))
  1381. const workspaces = await readWorkspaceBaseline(dropped.workspaceRemote)
  1382. if (!listed.result.ok) throw new Error('list failed')
  1383. expect(listed.result.value.items.some(item => item.sessionId === sessionId)).toBe(true)
  1384. expect(workspaces.items[0]?.sessionIds).toContain(sessionId)
  1385. await expect(dropped.sessions.create(req({
  1386. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1387. sessionId,
  1388. }))).resolves.toMatchObject({ result: { ok: true, value: { sessionId } } })
  1389. const rejecting = createFixtureApi({ empty: true, rejectPrompt: true })
  1390. const real = await rejecting.sessions.create(req({ sessionId: sid('fx-rejected') }))
  1391. if (!real.result.ok) throw new Error('session create failed')
  1392. const prompt = await rejecting.sessions.prompt(req({
  1393. sessionId: real.result.value.sessionId,
  1394. mode: 'queue' as const,
  1395. content: [{ type: 'text' as const, text: 'keep me' }],
  1396. }))
  1397. expect(prompt.result).toMatchObject({ ok: false, error: { code: 'session/agent-busy' } })
  1398. const imagePrompt = await rejecting.sessions.prompt(req({
  1399. sessionId: real.result.value.sessionId,
  1400. mode: 'queue' as const,
  1401. content: [{ type: 'image' as const, mediaType: 'image/png' as const, data: 'iVBORw0KGgo=' }],
  1402. }))
  1403. expect(imagePrompt.result).toMatchObject({
  1404. ok: false,
  1405. error: { code: 'session/attachment-invalid', details: { reason: 'IMAGE_DIMENSION_TOO_LARGE' } },
  1406. })
  1407. })
  1408. it('timing hooks: history delay + one-shot failure, silent append, and breakStreams end open generators', async () => {
  1409. const api = createFixtureApi()
  1410. const hooks = timing()
  1411. // One-shot transport failure after transit delay.
  1412. hooks.setHistoryDelay(5)
  1413. hooks.failNextHistory()
  1414. await expect(api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))).rejects.toThrow(/simulated history transport failure/)
  1415. hooks.setHistoryDelay(0)
  1416. // The failure was one-shot: the next call succeeds.
  1417. const ok = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))
  1418. expect(ok.result.ok).toBe(true)
  1419. // A durable append without a live frame creates a detectable seq gap.
  1420. const gapAbort = new AbortController()
  1421. const gapIterator = api.sessionRemote.follow(sid('fx-alpha'), gapAbort.signal)[Symbol.asyncIterator]()
  1422. const opening = await gapIterator.next()
  1423. if (opening.done || opening.value.type !== 'snapshot') throw new Error('follow opening snapshot missing')
  1424. hooks.appendSilent('fx-alpha', '静默丢帧')
  1425. hooks.appendUser('fx-alpha', '正常直播')
  1426. await expect(gapIterator.next()).rejects.toThrow(/stream skipped seq/)
  1427. // Reopening replaces the window with a complete snapshot containing both durable events.
  1428. const followAbort = new AbortController()
  1429. const controlAbort = new AbortController()
  1430. const followed: FixtureFollowFrame[] = []
  1431. const controlled: FixtureControlFrame[] = []
  1432. const following = (async () => {
  1433. for await (const frame of api.sessionRemote.follow(sid('fx-alpha'), followAbort.signal)) {
  1434. followed.push(frame)
  1435. }
  1436. })()
  1437. const controlling = (async () => {
  1438. for await (const frame of api.sessionRemote.control(controlAbort.signal)) controlled.push(frame)
  1439. })()
  1440. await new Promise(resolve => setTimeout(resolve, 10))
  1441. await vi.waitFor(() => {
  1442. const snapshot = followed.find(frame => frame.type === 'snapshot')
  1443. const events = snapshot === undefined ? [] : historyEvents(snapshot.records)
  1444. expect(events.some(event => JSON.stringify(event.data).includes('静默丢帧'))).toBe(true)
  1445. expect(events.some(event => JSON.stringify(event.data).includes('正常直播'))).toBe(true)
  1446. })
  1447. hooks.appendTitle('fx-alpha', 'Fixture 修订标题')
  1448. hooks.beginModelRetry('fx-alpha')
  1449. hooks.scheduleModelRetry('fx-alpha')
  1450. hooks.completeModelRetry('fx-alpha')
  1451. hooks.beginModelRetry('fx-alpha')
  1452. hooks.cancelModelRetryDuringBackoff('fx-alpha')
  1453. await vi.waitFor(() => {
  1454. expect(followed.some(frame => frame.type === 'event' && (frame.event as { type: string }).type === 'llm/retry')).toBe(true)
  1455. expect(followed.some(frame => frame.type === 'event' && JSON.stringify(frame.event.data).includes('重试后的完整回复'))).toBe(true)
  1456. expect(followed.some(frame => frame.type === 'event'
  1457. && frame.event.type === 'turn/end'
  1458. && frame.event.data.reason.kind === 'aborted')).toBe(true)
  1459. expect(controlled.some(frame => frame.type === 'projection'
  1460. && frame.key === 'title'
  1461. && frame.value === 'Fixture 修订标题')).toBe(true)
  1462. })
  1463. expect(followed.some(frame => frame.type === 'event' && (frame.event as { type: string }).type === 'session/title')).toBe(true)
  1464. // Paging and resumed follow agree on the recovered durable event.
  1465. const repull = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))
  1466. if (!repull.result.ok) throw new Error('repull failed')
  1467. expect(JSON.stringify(repull.result.value.records)).toContain('静默丢帧')
  1468. // breakStreams force-ends follow and control without client aborts.
  1469. await new Promise(resolve => setTimeout(resolve, 10))
  1470. hooks.breakStreams()
  1471. await following
  1472. await controlling
  1473. expect(followAbort.signal.aborted).toBe(false)
  1474. expect(controlAbort.signal.aborted).toBe(false)
  1475. })
  1476. it('paces the opt-in reasoning stress hook from an external interval', async () => {
  1477. vi.useFakeTimers()
  1478. vi.setSystemTime(0)
  1479. const api = createFixtureApi()
  1480. const hooks = timing()
  1481. expect(hooks.reasoningChunkStormState()).toBeNull()
  1482. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 0, 1, 16)).toThrow(/chunk count/)
  1483. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 1, 0, 16)).toThrow(/chunks per interval/)
  1484. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 1, 1, 0)).toThrow(/reasoning interval/)
  1485. const abort = new AbortController()
  1486. try {
  1487. const streamed = collectValues(api.sessionRemote.follow(sid('fx-alpha'), abort.signal), abort, frames => frames.some(frame => (
  1488. frame.type === 'assistant-stream'
  1489. && frame.frame.type === 'chunk'
  1490. && isReasoningDeltaChunk(frame.frame.chunk)
  1491. && frame.frame.chunk.text.includes('REASONING_STRESS_COMPLETE')
  1492. )))
  1493. const marker = hooks.startReasoningChunkStorm('fx-alpha', 3, 2, 16)
  1494. expect(() => hooks.startReasoningChunkStorm('fx-alpha', 1, 1, 16)).toThrow(/already running/)
  1495. expect(hooks.reasoningChunkStormState()).toMatchObject({ emitted: 0, emitting: true, marker })
  1496. await vi.advanceTimersByTimeAsync(0)
  1497. expect(hooks.reasoningChunkStormState()).toMatchObject({ emitted: 2, emitting: true })
  1498. await vi.advanceTimersByTimeAsync(16)
  1499. expect(hooks.reasoningChunkStormState()).toEqual({
  1500. sessionId: 'fx-alpha', chunkCount: 3, chunksPerInterval: 2, intervalMs: 16,
  1501. emitted: 3, marker, emitting: false,
  1502. })
  1503. const frames = await streamed
  1504. const deltas = frames.flatMap(frame => (
  1505. frame.type === 'assistant-stream'
  1506. && frame.frame.type === 'chunk'
  1507. && isReasoningDeltaChunk(frame.frame.chunk)
  1508. ? [frame.frame.chunk.text]
  1509. : []
  1510. ))
  1511. expect(deltas).toEqual(['推理', '推理', `\n${marker}`])
  1512. } finally {
  1513. abort.abort()
  1514. vi.useRealTimers()
  1515. }
  1516. })
  1517. })
  1518. describe('fixture Connection RPC', () => {
  1519. afterEach(() => {
  1520. vi.restoreAllMocks()
  1521. vi.unstubAllGlobals()
  1522. })
  1523. it('covers the migrated Remote dispatch table', async () => {
  1524. const rpc = createFixtureConnectionRpc()
  1525. const sessions = createSessionClient(rpc)
  1526. const workspaces = createWorkspaceClient(rpc)
  1527. expect((await sessions.search(
  1528. { query: 'fixture' },
  1529. new AbortController().signal,
  1530. )).result.ok).toBe(true)
  1531. const created = await sessions.create({})
  1532. if (!created.result.ok) throw new Error('create failed')
  1533. const id = created.result.value.sessionId
  1534. expect((await sessions.history({ sessionId: id })).result.ok).toBe(true)
  1535. expect((await sessions.prompt({ sessionId: id, mode: 'queue', content: [{ type: 'text', text: '嗨' }] })).result.ok).toBe(true)
  1536. expect((await sessions.cancel({ sessionId: id })).result.ok).toBe(true)
  1537. expect((await readWorkspaceBaseline(createWorkspaceRemote(rpc))).items).not.toHaveLength(0)
  1538. const workspace = await workspaces.create({ path: '/tmp/fixture-workspaces/via-client' })
  1539. if (!workspace.result.ok) throw new Error('workspace create failed')
  1540. expect(workspace.result.value.workspace.title).toBe('via-client')
  1541. const wsid = workspace.result.value.workspace.workspaceId
  1542. const renamed = await workspaces.rename({ workspaceId: wsid, title: 'via-client-2' })
  1543. if (!renamed.result.ok) throw new Error('workspace rename failed')
  1544. expect(renamed.result.value.workspace.title).toBe('via-client-2')
  1545. const attached = await sessions.create({ workspaceId: wsid })
  1546. if (!attached.result.ok) throw new Error('attached create failed')
  1547. const moved = await workspaces.insertSessionBefore({ workspaceId: wsid, sessionId: attached.result.value.sessionId })
  1548. if (!moved.result.ok) throw new Error('workspace move failed')
  1549. expect(moved.result.value.workspace.sessionIds).toEqual([attached.result.value.sessionId])
  1550. })
  1551. it('folds the goal lifecycle over the Goal Remotes', async () => {
  1552. const rpc = createFixtureConnectionRpc()
  1553. const sessions = createSessionClient(rpc)
  1554. const created = await sessions.create({})
  1555. if (!created.result.ok) throw new Error('create failed')
  1556. const id = created.result.value.sessionId
  1557. const goal = (endpoint: string, args: Record<string, unknown>) =>
  1558. rpc.call('/api', endpoint, { args: { agentId: id, ...args } })
  1559. // create → edit → pause → resume → complete → clear; each mutation advances the CAS
  1560. // revision by one (state rides the projection frames).
  1561. const goalCreated = await goal('goals/create', { request: { objective: 'ship it' } })
  1562. if (!goalCreated.ok) throw new Error('goal create failed')
  1563. const { id: goalId, revision } = (goalCreated.value as { ref: { id: string; revision: number } }).ref
  1564. expect(revision).toBe(1)
  1565. const ref = (at: number) => ({ id: goalId, revision: at })
  1566. expect((await goal('goals/edit', { ref: ref(1), request: { objective: 'ship it v2' } })).ok).toBe(true)
  1567. expect(await goal('goals/get', {})).toMatchObject({
  1568. ok: true, value: { objective: 'ship it v2', revision: 2, activation: 'armed' },
  1569. })
  1570. timing().disarmOnlyGoal()
  1571. expect(await goal('goals/get', {})).toMatchObject({
  1572. ok: true, value: { objective: 'ship it v2', revision: 2, activation: 'disarmed' },
  1573. })
  1574. expect((await goal('goals/pause', { ref: ref(2) })).ok).toBe(true)
  1575. expect((await goal('goals/resume', { ref: ref(3) })).ok).toBe(true)
  1576. // A stale ref loses the CAS check.
  1577. expect((await goal('goals/pause', { ref: ref(1) })).ok).toBe(false)
  1578. expect((await goal('goals/complete', { ref: ref(4) })).ok).toBe(true)
  1579. // complete → complete is an invalid transition.
  1580. expect((await goal('goals/complete', { ref: ref(5) })).ok).toBe(false)
  1581. expect(await goal('goals/clear', { ref: ref(5) })).toEqual({ ok: true, value: ref(6) })
  1582. const goalHistory = await sessions.history({ sessionId: id })
  1583. if (!goalHistory.result.ok) throw new Error('goal history failed')
  1584. const goalEvents = historyEvents(goalHistory.result.value.records).map(event => event as unknown as {
  1585. type: string
  1586. data: {
  1587. operation?: string
  1588. source?: { kind?: string; round?: number }
  1589. }
  1590. })
  1591. const goalChanges = goalEvents.filter(event => event.type === 'goal/change')
  1592. expect(goalChanges.map(event => event.data.operation))
  1593. .toEqual(['create', 'edit', 'pause', 'resume', 'complete', 'clear'])
  1594. expect(goalEvents.some(event => event.type === 'user/message'
  1595. && event.data.source?.kind === 'goal' && event.data.source.round === 0)).toBe(false)
  1596. })
  1597. it('maps empty, prompt-reject, and workspace-first query scenarios', async () => {
  1598. vi.stubGlobal('location', {
  1599. search: '?fixture=empty&fixturePrompt=reject&fixtureFrames=workspace-first',
  1600. })
  1601. const rpc = createFixtureConnectionRpc()
  1602. const sessions = createSessionClient(rpc)
  1603. const workspaces = createWorkspaceClient(rpc)
  1604. const workspaceRemote = createWorkspaceRemote(rpc)
  1605. await expect(sessions.list({})).resolves.toMatchObject({ result: { ok: true, value: { items: [] } } })
  1606. const made = await workspaces.create({ path: '/tmp/fixture-workspaces/query-workspace' })
  1607. if (!made.result.ok) throw new Error('workspace create failed')
  1608. const hostAbort = new AbortController()
  1609. const workspaceAbort = new AbortController()
  1610. const hostFrames = collectValues(
  1611. openFixtureRemoteEvents(rpc, hostAbort.signal),
  1612. hostAbort,
  1613. frames => frames.length === 1,
  1614. )
  1615. const workspaceFrames = collectValues(
  1616. workspaceRemote.follow(workspaceAbort.signal),
  1617. workspaceAbort,
  1618. frames => frames.some(frame => frame.type === 'upsert'
  1619. && frame.workspace.sessionIds.includes(sid('fx-query-session'))),
  1620. )
  1621. await new Promise(resolve => setTimeout(resolve, 10))
  1622. const sessionId = sid('fx-query-session')
  1623. const created = await sessions.create({
  1624. workspaceId: made.result.value.workspace.workspaceId,
  1625. sessionId,
  1626. })
  1627. expect(created.result).toMatchObject({ ok: true, value: { sessionId } })
  1628. expect((await workspaceFrames).at(-1)).toMatchObject({
  1629. type: 'upsert',
  1630. workspace: { sessionIds: [sessionId] },
  1631. })
  1632. expect((await hostFrames)[0]).toMatchObject({
  1633. event: 'api-session/added',
  1634. })
  1635. const rejected = await sessions.prompt({
  1636. sessionId,
  1637. mode: 'queue',
  1638. content: [{ type: 'text', text: 'retain' }],
  1639. })
  1640. expect(rejected.result).toMatchObject({ ok: false, error: { code: 'session/agent-busy' } })
  1641. })
  1642. it('maps attach-failure and dropped-response query scenarios', async () => {
  1643. vi.stubGlobal('location', { search: '?fixture&fixtureAttach=fail' })
  1644. const partial = createFixtureConnectionRpc()
  1645. const partialResult = await createSessionClient(partial).create({
  1646. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1647. sessionId: sid('fx-query-partial'),
  1648. })
  1649. expect(partialResult.result).toMatchObject({
  1650. ok: false,
  1651. error: { code: 'session/workspace-attach-failed', details: { sessionId: 'fx-query-partial' } },
  1652. })
  1653. vi.stubGlobal('location', { search: '?fixture&fixtureSessionCreate=drop-response' })
  1654. const dropped = createFixtureConnectionRpc()
  1655. await expect(createSessionClient(dropped).create({
  1656. workspaceId: 'fx-ws-fixture' as WorkspaceId,
  1657. sessionId: sid('fx-query-dropped'),
  1658. })).rejects.toThrow(/dropped session\.create response/)
  1659. })
  1660. })