fetch-carrier.spec.ts 37 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816
  1. import { CommandId } from '@deepseek-ai/dsh-commands/brand'
  2. import { describe, expect, it, vi } from 'vitest'
  3. import type { ApiProxy, HostFrame, MuxFrame } from '../src/api/index.ts'
  4. import type { ClientResponse, RpcMessage, RpcReceipt, RpcRequest } from '../src/api/rpc.ts'
  5. import { RpcId } from '../src/api/rpc.ts'
  6. import { toFetchHandler } from '../src/fetch/handler.ts'
  7. import { AbstractApiClient, InProcessApiClient } from '../src/fetch/client.ts'
  8. /** Minimal in-memory ApiProxy: echoes rpcIds, scripts one frame per stream. */
  9. function fakeApi(overrides: Partial<{ muxFrames: MuxFrame[]; hostFrames: HostFrame[]; crashOn: string }> = {}): ApiProxy {
  10. const muxFrames = overrides.muxFrames ?? [{ type: 'session/subscribed', sessionId: 's1' as never, lastSeq: -1 }]
  11. const hostFrames = overrides.hostFrames ?? [{ type: 'host/session-removed', sessionId: 's1' as never }]
  12. async function * stream<F>(frames: F[], signal: AbortSignal): AsyncGenerator<RpcRequest<F>> {
  13. for (const payload of frames) {
  14. if (signal.aborted) return
  15. yield { rpcId: RpcId(`frame-${String(frames.indexOf(payload))}`), payload }
  16. }
  17. }
  18. return {
  19. sessions: {
  20. async list(request) {
  21. if (overrides.crashOn === 'session.list') throw new Error('impl crashed')
  22. return { rpcId: request.rpcId, result: { ok: true, value: { items: [] } } }
  23. },
  24. async search(request, signal) {
  25. if (request.payload.query === 'hang') {
  26. if (!signal.aborted) {
  27. await new Promise<void>((resolve) => {
  28. signal.addEventListener('abort', () => { resolve() }, { once: true })
  29. })
  30. }
  31. return {
  32. rpcId: request.rpcId,
  33. result: { ok: false, error: { code: 'cancelled', message: 'aborted', details: {} } },
  34. }
  35. }
  36. return {
  37. rpcId: request.rpcId,
  38. result: {
  39. ok: true,
  40. value: { items: [{ sessionId: 's1' as never, snippet: 'fixture match' }], hasMore: false },
  41. },
  42. }
  43. },
  44. async create(request) {
  45. return { rpcId: request.rpcId, result: { ok: true, value: { sessionId: 's-new' as never } } }
  46. },
  47. async history(request) {
  48. if (request.payload.sessionId === ('with-projections' as never)) {
  49. return {
  50. rpcId: request.rpcId,
  51. result: { ok: true, value: { events: [], hasMore: false, projections: { asOfSeq: 9, values: { todos: [{ content: 'current', status: 'in_progress' as const }] } } } },
  52. }
  53. }
  54. return {
  55. rpcId: request.rpcId,
  56. result: { ok: false, error: { code: 'session-not-found', message: 'nope', details: { sessionId: request.payload.sessionId } } },
  57. }
  58. },
  59. async models(request) {
  60. return {
  61. rpcId: request.rpcId,
  62. result: {
  63. ok: true,
  64. value: {
  65. current: { provider: 'deepseek-official', model: 'deepseek-v4-flash' },
  66. groups: [],
  67. failures: [],
  68. },
  69. },
  70. }
  71. },
  72. async selectModel(request) {
  73. return {
  74. rpcId: request.rpcId,
  75. result: {
  76. ok: true,
  77. value: {
  78. selected: {
  79. provider: request.payload.provider,
  80. model: request.payload.model,
  81. ...request.payload.reasoningEffort === undefined
  82. ? {}
  83. : { reasoningEffort: request.payload.reasoningEffort },
  84. },
  85. },
  86. },
  87. }
  88. },
  89. async rename(request) {
  90. return { rpcId: request.rpcId, result: { ok: true, value: { title: request.payload.title, seq: 0 } } }
  91. },
  92. async fork(request) {
  93. return { rpcId: request.rpcId, result: { ok: true, value: { sessionId: 's-fork' as never } } }
  94. },
  95. async prompt(request) {
  96. return { rpcId: request.rpcId, result: { ok: true, value: { accepted: true as const } } }
  97. },
  98. async updateQueue(request) {
  99. return { rpcId: request.rpcId, result: { ok: true, value: { accepted: true as const } } }
  100. },
  101. async cancel(request) {
  102. return { rpcId: request.rpcId, result: { ok: true, value: { accepted: true as const } } }
  103. },
  104. },
  105. subagents: {
  106. async list(request) {
  107. return { rpcId: request.rpcId, result: { ok: true, value: { entries: [], parentAvailable: false } } }
  108. },
  109. async history(request) {
  110. return { rpcId: request.rpcId, result: { ok: true, value: { events: [], hasMore: false } } }
  111. },
  112. async prompt(request, signal) {
  113. if (request.payload.content.some(block => block.type === 'text' && block.text === 'hang')) {
  114. if (!signal.aborted) {
  115. await new Promise<void>((resolve) => {
  116. signal.addEventListener('abort', () => { resolve() }, { once: true })
  117. })
  118. }
  119. return {
  120. rpcId: request.rpcId,
  121. result: { ok: false, error: { code: 'cancelled' as const, message: 'aborted', details: {} } },
  122. }
  123. }
  124. return {
  125. rpcId: request.rpcId,
  126. result: { ok: true, value: { messageId: 'message-1' as never } },
  127. }
  128. },
  129. },
  130. host: {
  131. async describe(request) {
  132. return { rpcId: request.rpcId, result: { ok: true, value: { version: 'v', cwd: '/w', attachedSessions: 0 } } }
  133. },
  134. async pickDirectory(request) {
  135. return { rpcId: request.rpcId, result: { ok: true, value: { path: null } } }
  136. },
  137. async listDirectory(request) {
  138. return { rpcId: request.rpcId, result: { ok: true, value: { path: '/w', home: '/w', crumbs: [{ name: '/', path: '/', hidden: false }], entries: [], truncated: false } } }
  139. },
  140. async createDirectory(request) {
  141. return { rpcId: request.rpcId, result: { ok: true, value: { path: '/w/new' } } }
  142. },
  143. async openPath(request) {
  144. return { rpcId: request.rpcId, result: { ok: true, value: { opened: true as const } } }
  145. },
  146. },
  147. workspace: {
  148. async list(request) {
  149. return { rpcId: request.rpcId, result: { ok: true, value: { items: [], archivedSessionIds: [] } } }
  150. },
  151. async create(request) {
  152. return {
  153. rpcId: request.rpcId,
  154. result: { ok: true, value: { workspace: { workspaceId: 'w1' as never, path: '/w', title: 'w', sessionIds: [], createdAt: 't', updatedAt: 't' }, created: true } },
  155. }
  156. },
  157. async rename(request) {
  158. return {
  159. rpcId: request.rpcId,
  160. result: { ok: true, value: { workspace: { workspaceId: 'w1' as never, path: '/w', title: 'w', sessionIds: [], createdAt: 't', updatedAt: 't' } } },
  161. }
  162. },
  163. async delete(request) {
  164. return { rpcId: request.rpcId, result: { ok: true, value: { deleted: true as const } } }
  165. },
  166. async insertSessionBefore(request) {
  167. return {
  168. rpcId: request.rpcId,
  169. result: { ok: true, value: { workspace: { workspaceId: 'w1' as never, path: '/w', title: 'w', sessionIds: [], createdAt: 't', updatedAt: 't' } } },
  170. }
  171. },
  172. async archiveSession(request) {
  173. return { rpcId: request.rpcId, result: { ok: true, value: { archivedSessionIds: [request.payload.sessionId] } } }
  174. },
  175. },
  176. commands: {
  177. async list(request) {
  178. return { rpcId: request.rpcId, result: { ok: true, value: { commands: [{ name: 'plan', description: 'Toggle plan mode', input: { hint: 'on|off' } }] } } }
  179. },
  180. async execute(request, signal) {
  181. if (request.payload.line === '/hang') {
  182. // Cooperative hang: settles only through the carrier signal (sticky
  183. // abort checked first — listeners never fire retroactively).
  184. if (!signal.aborted) {
  185. await new Promise<void>((resolve) => { signal.addEventListener('abort', () => { resolve() }, { once: true }) })
  186. }
  187. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'cancelled', message: 'aborted', details: {} } } }
  188. }
  189. if (request.payload.line.startsWith('/plan')) {
  190. return { rpcId: request.rpcId, result: { ok: true, value: { matched: true, commandId: CommandId('cmd-x') } } }
  191. }
  192. return { rpcId: request.rpcId, result: { ok: true, value: { matched: false } } }
  193. },
  194. },
  195. agentPresets: {
  196. list(request: RpcRequest<{}>) {
  197. return Promise.resolve({
  198. rpcId: request.rpcId,
  199. result: { ok: true as const, value: { presets: [], authorable: false } },
  200. })
  201. },
  202. select(request: RpcRequest<{ agentPreset: string }>) {
  203. const value = { agentPreset: request.payload.agentPreset }
  204. return Promise.resolve({ rpcId: request.rpcId, result: { ok: true as const, value } })
  205. },
  206. read(request: RpcRequest<{ agentPreset: string }>) {
  207. const value = { agentPreset: request.payload.agentPreset, trust: 'user' as const, content: '', writable: true }
  208. return Promise.resolve({ rpcId: request.rpcId, result: { ok: true as const, value } })
  209. },
  210. write(request: RpcRequest<{ agentPreset: string }>) {
  211. const value = { agentPreset: request.payload.agentPreset }
  212. return Promise.resolve({ rpcId: request.rpcId, result: { ok: true as const, value } })
  213. },
  214. remove(request: RpcRequest<{ agentPreset: string }>) {
  215. return Promise.resolve({ rpcId: request.rpcId, result: { ok: true as const, value: {} } })
  216. },
  217. },
  218. skills: {
  219. async list(request) {
  220. return { rpcId: request.rpcId, result: { ok: true, value: { skills: [{ name: 'commit-helper', description: 'Git commits' }] } } }
  221. },
  222. },
  223. goals: {
  224. async create(request) {
  225. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  226. },
  227. async edit(request) {
  228. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  229. },
  230. async pause(request) {
  231. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  232. },
  233. async resume(request) {
  234. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  235. },
  236. async complete(request) {
  237. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  238. },
  239. async clear(request) {
  240. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  241. },
  242. },
  243. settings: {
  244. async describe(request) {
  245. return { rpcId: request.rpcId, result: { ok: true, value: { writable: true, hasDocument: false, namespaces: [] } } }
  246. },
  247. async openDocument(request) {
  248. return { rpcId: request.rpcId, result: { ok: true, value: { opened: true as const } } }
  249. },
  250. async update(request) {
  251. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'settings-rejected', message: 'stub', details: { ns: request.payload.ns } } } }
  252. },
  253. async replace(request) {
  254. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'settings-rejected', message: 'stub', details: { ns: request.payload.ns } } } }
  255. },
  256. async mutate(request) {
  257. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'settings-rejected', message: 'stub', details: { ns: request.payload.ns } } } }
  258. },
  259. },
  260. credentials: {
  261. async describe(request) {
  262. return { rpcId: request.rpcId, result: { ok: true, value: { credentials: {} } } }
  263. },
  264. async set(request) {
  265. return { rpcId: request.rpcId, result: { ok: true, value: {} } }
  266. },
  267. async unset(request) {
  268. return { rpcId: request.rpcId, result: { ok: true, value: {} } }
  269. },
  270. },
  271. llm: {
  272. async providers(request) {
  273. return { rpcId: request.rpcId, result: { ok: true, value: { providers: [] } } }
  274. },
  275. async models(request) {
  276. return { rpcId: request.rpcId, result: { ok: true, value: { groups: [], failures: [] } } }
  277. },
  278. async discoverModels(request) {
  279. return { rpcId: request.rpcId, result: { ok: true, value: { models: [] } } }
  280. },
  281. },
  282. events: {
  283. mux: (_request, signal) => stream(muxFrames, signal),
  284. host: (_request, signal) => stream(hostFrames, signal),
  285. },
  286. async respond(message: ClientResponse): Promise<RpcReceipt> {
  287. return message.rpcId === 'known' ? { accepted: true } : { accepted: false, reason: 'not-pending' }
  288. },
  289. }
  290. }
  291. function client(api: ApiProxy = fakeApi(), timeoutMs?: number): InProcessApiClient {
  292. return new InProcessApiClient(toFetchHandler(api), timeoutMs)
  293. }
  294. async function collect<F>(stream: AsyncIterable<RpcRequest<F>>): Promise<RpcRequest<F>[]> {
  295. const out: RpcRequest<F>[] = []
  296. for await (const envelope of stream) out.push(envelope)
  297. return out
  298. }
  299. describe('unary round trip (handler ⇄ client, no network)', () => {
  300. it('carries a success result and echoes the minted rpcId', async () => {
  301. const response = await client().sessions.list({})
  302. expect(response.result).toEqual({ ok: true, value: { items: [] } })
  303. expect(response.rpcId).toMatch(/[0-9a-f-]{36}/)
  304. })
  305. it('carries the tail-page projections block through the wire schema (Zod must not strip it)', async () => {
  306. const response = await client().sessions.history({ sessionId: 'with-projections' as never })
  307. expect(response.result.ok).toBe(true)
  308. if (response.result.ok) {
  309. expect(response.result.value.projections).toEqual(
  310. { asOfSeq: 9, values: { todos: [{ content: 'current', status: 'in_progress' }] } },
  311. )
  312. }
  313. })
  314. it('carries a business error as 200 + error result', async () => {
  315. const response = await client().sessions.history({ sessionId: 'missing' as never })
  316. expect(response.result.ok).toBe(false)
  317. if (!response.result.ok) expect(response.result.error.code).toBe('session-not-found')
  318. })
  319. it('covers create/prompt/updateQueue/cancel/describe passthrough', async () => {
  320. const c = client()
  321. expect((await c.sessions.search({ query: 'fixture' })).result).toEqual({
  322. ok: true,
  323. value: { items: [{ sessionId: 's1', snippet: 'fixture match' }], hasMore: false },
  324. })
  325. expect((await c.sessions.create({})).result.ok).toBe(true)
  326. expect((await c.sessions.models({ sessionId: 's' as never })).result.ok).toBe(true)
  327. const selected = await c.sessions.selectModel({
  328. sessionId: 's' as never,
  329. provider: 'deepseek-official',
  330. model: 'deepseek-v4-flash',
  331. reasoningEffort: 'max',
  332. })
  333. expect(selected.result).toMatchObject({
  334. ok: true,
  335. value: {
  336. selected: {
  337. provider: 'deepseek-official',
  338. model: 'deepseek-v4-flash',
  339. reasoningEffort: 'max',
  340. },
  341. },
  342. })
  343. const renamed = await c.sessions.rename({ sessionId: 's' as never, title: 'named' })
  344. expect(renamed.result).toMatchObject({ ok: true, value: { title: 'named', seq: 0 } })
  345. expect((await c.sessions.prompt({ sessionId: 's' as never, mode: 'queue', content: [{ type: 'text', text: 'x' }] })).result.ok).toBe(true)
  346. expect((await c.sessions.updateQueue({
  347. sessionId: 's' as never,
  348. itemId: 'item-1' as never,
  349. action: { kind: 'remove' },
  350. })).result.ok).toBe(true)
  351. expect((await c.sessions.cancel({ sessionId: 's' as never })).result.ok).toBe(true)
  352. expect((await c.host.describe({})).result.ok).toBe(true)
  353. })
  354. it('round-trips every agent-preset method, authoring included', async () => {
  355. const c = client()
  356. // The whole domain crosses the carrier: the roster a picker reads, the
  357. // per-session switch, and the three authoring calls the settings editor
  358. // makes. Each has its own request schema, so a registration missing from
  359. // either half fails here rather than in the browser.
  360. expect((await c.agentPresets.list({})).result).toEqual({
  361. ok: true, value: { presets: [], authorable: false },
  362. })
  363. expect((await c.agentPresets.select({ sessionId: 's' as never, agentPreset: 'minimal' })).result)
  364. .toEqual({ ok: true, value: { agentPreset: 'minimal' } })
  365. expect((await c.agentPresets.read({ agentPreset: 'mine' })).result).toEqual({
  366. ok: true, value: { agentPreset: 'mine', trust: 'user', content: '', writable: true },
  367. })
  368. expect((await c.agentPresets.write({ agentPreset: 'mine', content: '- id: x\n' })).result)
  369. .toEqual({ ok: true, value: { agentPreset: 'mine' } })
  370. expect((await c.agentPresets.remove({ agentPreset: 'mine' })).result).toEqual({ ok: true, value: {} })
  371. })
  372. it('round-trips the native picker without the default unary timeout', async () => {
  373. const api = fakeApi()
  374. api.host.pickDirectory = async (request) => {
  375. await new Promise(resolve => setTimeout(resolve, 15))
  376. return { rpcId: request.rpcId, result: { ok: true, value: { path: '/tmp/project' } } }
  377. }
  378. const response = await client(api, 1).host.pickDirectory({})
  379. expect(response.result).toEqual({ ok: true, value: { path: '/tmp/project' } })
  380. })
  381. it('round-trips the browse listing and creation calls through the wire form', async () => {
  382. const c = client()
  383. const listed = await c.host.listDirectory({ path: '/w' })
  384. expect(listed.result).toEqual({
  385. ok: true,
  386. value: { path: '/w', home: '/w', crumbs: [{ name: '/', path: '/', hidden: false }], entries: [], truncated: false },
  387. })
  388. const home = await c.host.listDirectory({})
  389. expect(home.result).toMatchObject({ ok: true, value: { home: '/w' } })
  390. const created = await c.host.createDirectory({ path: '/w', name: 'fresh' })
  391. expect(created.result).toEqual({ ok: true, value: { path: '/w/new' } })
  392. })
  393. it('round-trips host.openPath through the wire form', async () => {
  394. const api = fakeApi()
  395. let opened: string | undefined
  396. api.host.openPath = async (request) => {
  397. opened = request.payload.path
  398. return { rpcId: request.rpcId, result: { ok: true, value: { opened: true as const } } }
  399. }
  400. const response = await client(api).host.openPath({ path: '/tmp/a.txt' })
  401. expect(opened).toBe('/tmp/a.txt')
  402. expect(response.result).toEqual({ ok: true, value: { opened: true } })
  403. })
  404. it('round-trips command.list / command.execute / skill.list through the wire form', async () => {
  405. const c = client()
  406. const list = await c.commands.list({ sessionId: 's' as never })
  407. expect(list.result).toEqual({ ok: true, value: { commands: [{ name: 'plan', description: 'Toggle plan mode', input: { hint: 'on|off' } }] } })
  408. const hit = await c.commands.execute({ sessionId: 's' as never, line: '/plan off' })
  409. expect(hit.result).toEqual({ ok: true, value: { matched: true, commandId: 'cmd-x' } })
  410. const miss = await c.commands.execute({ sessionId: 's' as never, line: '/nope' })
  411. expect(miss.result).toEqual({ ok: true, value: { matched: false } })
  412. const skills = await c.skills.list({ sessionId: 's' as never })
  413. expect(skills.result).toEqual({ ok: true, value: { skills: [{ name: 'commit-helper', description: 'Git commits' }] } })
  414. })
  415. it('lets command.execute finish after the 30-second default unary deadline', async () => {
  416. vi.useFakeTimers()
  417. const timeoutSpy = vi.spyOn(AbortSignal, 'timeout').mockImplementation((milliseconds) => {
  418. const controller = new AbortController()
  419. setTimeout(() => {
  420. controller.abort(new DOMException('The operation was aborted due to timeout', 'TimeoutError'))
  421. }, milliseconds)
  422. return controller.signal
  423. })
  424. try {
  425. const api = fakeApi()
  426. api.commands.execute = async (request) => {
  427. await new Promise(resolve => setTimeout(resolve, 30_001))
  428. return {
  429. rpcId: request.rpcId,
  430. result: { ok: true, value: { matched: true, commandId: CommandId('cmd-slow') } },
  431. }
  432. }
  433. const execution = client(api).commands.execute({ sessionId: 's' as never, line: '/slow' })
  434. const assertion = expect(execution).resolves.toMatchObject({
  435. result: { ok: true, value: { matched: true, commandId: 'cmd-slow' } },
  436. })
  437. await Promise.all([
  438. vi.advanceTimersByTimeAsync(30_001),
  439. assertion,
  440. ])
  441. expect(timeoutSpy).not.toHaveBeenCalled()
  442. } finally {
  443. timeoutSpy.mockRestore()
  444. vi.useRealTimers()
  445. }
  446. })
  447. it('round-trips the subagent domain through the wire form', async () => {
  448. const c = client()
  449. expect((await c.subagents.list({ parentSessionId: 'parent' as never })).result)
  450. .toEqual({ ok: true, value: { entries: [], parentAvailable: false } })
  451. expect((await c.subagents.history({
  452. parentSessionId: 'parent' as never,
  453. childSessionId: 'child' as never,
  454. mode: 'one-shot',
  455. })).result).toEqual({ ok: true, value: { events: [], hasMore: false } })
  456. expect((await c.subagents.prompt({
  457. parentSessionId: 'parent' as never,
  458. childSessionId: 'child' as never,
  459. mode: 'continuable',
  460. content: [],
  461. })).result).toEqual({ ok: true, value: { messageId: 'message-1' } })
  462. })
  463. it('keeps caller and connection aborts on command.execute', async () => {
  464. const api = fakeApi()
  465. const started = Promise.withResolvers<AbortSignal>()
  466. api.commands.execute = async (request, signal) => {
  467. started.resolve(signal)
  468. if (!signal.aborted) {
  469. await new Promise<void>((resolve) => {
  470. signal.addEventListener('abort', () => { resolve() }, { once: true })
  471. })
  472. }
  473. return {
  474. rpcId: request.rpcId,
  475. result: { ok: false, error: { code: 'cancelled', message: 'aborted', details: {} } },
  476. }
  477. }
  478. const controller = new AbortController()
  479. const execution = client(api).commands.execute(
  480. { sessionId: 's' as never, line: '/hang' },
  481. controller.signal,
  482. )
  483. const handlerSignal = await started.promise
  484. controller.abort(new Error('connection closed'))
  485. await expect(execution).rejects.toThrow('connection closed')
  486. expect(handlerSignal.aborted).toBe(true)
  487. })
  488. it('propagates the carrier Request signal into command.execute', async () => {
  489. const handler = toFetchHandler(fakeApi())
  490. const controller = new AbortController()
  491. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-sig', method: 'command.execute', payload: { sessionId: 's', line: '/hang' } })
  492. // The fake's /hang settles only when the invoke-level signal aborts: a
  493. // completed response with the cancelled error proves req.signal reached it.
  494. const pending = handler.fetch(new Request('http://x/api/command.execute', { method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal }))
  495. controller.abort()
  496. const response = await pending
  497. const parsed = await response.json() as { rpcId: string; result: { ok: boolean; error?: { code: string } } }
  498. expect(parsed.rpcId).toBe('r-sig')
  499. expect(parsed.result.error?.code).toBe('cancelled')
  500. })
  501. it('propagates the carrier Request signal into session.search', async () => {
  502. const handler = toFetchHandler(fakeApi())
  503. const controller = new AbortController()
  504. const body = JSON.stringify({
  505. type: 'client-request',
  506. rpcId: 'r-search-sig',
  507. method: 'session.search',
  508. payload: { query: 'hang' },
  509. })
  510. const pending = handler.fetch(new Request(
  511. 'http://x/api/session.search',
  512. { method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal },
  513. ))
  514. controller.abort()
  515. const response = await pending
  516. const parsed = await response.json() as {
  517. rpcId: string
  518. result: { error?: { code: string } }
  519. }
  520. expect(parsed.rpcId).toBe('r-search-sig')
  521. expect(parsed.result.error?.code).toBe('cancelled')
  522. })
  523. it('propagates the carrier Request signal into subagent.prompt', async () => {
  524. const handler = toFetchHandler(fakeApi())
  525. const controller = new AbortController()
  526. const body = JSON.stringify({
  527. type: 'client-request',
  528. rpcId: 'r-subagent-sig',
  529. method: 'subagent.prompt',
  530. payload: {
  531. parentSessionId: 'parent',
  532. childSessionId: 'child',
  533. mode: 'continuable',
  534. content: [{ type: 'text', text: 'hang' }],
  535. },
  536. })
  537. const pending = handler.fetch(new Request(
  538. 'http://x/api/subagent.prompt',
  539. { method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal },
  540. ))
  541. controller.abort()
  542. const response = await pending
  543. const parsed = await response.json() as {
  544. rpcId: string
  545. result: { error?: { code: string } }
  546. }
  547. expect(parsed.rpcId).toBe('r-subagent-sig')
  548. expect(parsed.result.error?.code).toBe('cancelled')
  549. })
  550. it('propagates the carrier Request signal into host.pickDirectory', async () => {
  551. const api = fakeApi()
  552. api.host.pickDirectory = async (request, signal) => {
  553. if (!signal.aborted) {
  554. await new Promise<void>((resolve) => {
  555. signal.addEventListener('abort', () => { resolve() }, { once: true })
  556. })
  557. }
  558. return {
  559. rpcId: request.rpcId,
  560. result: { ok: false, error: { code: 'cancelled', message: 'aborted', details: {} } },
  561. }
  562. }
  563. const handler = toFetchHandler(api)
  564. const controller = new AbortController()
  565. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-picker', method: 'host.pickDirectory', payload: {} })
  566. const pending = handler.fetch(new Request('http://x/api/host.pickDirectory', {
  567. method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal,
  568. }))
  569. controller.abort()
  570. const parsed = await (await pending).json() as { result: { error?: { code: string } } }
  571. expect(parsed.result.error?.code).toBe('cancelled')
  572. })
  573. })
  574. describe('handler carrier-layer statuses', () => {
  575. const handler = toFetchHandler(fakeApi())
  576. it('404s unknown paths and non-POST non-stream methods', async () => {
  577. expect((await handler.fetch(new Request('http://x/other', { method: 'POST', headers: { 'content-type': 'application/json' }, body: '{}' }))).status).toBe(404)
  578. expect((await handler.fetch(new Request('http://x/api/session.list', { method: 'GET' }))).status).toBe(404)
  579. expect((await handler.fetch(new Request('http://x/api/no.such', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ type: 'client-request', rpcId: 'r', method: 'no.such', payload: {} }) }))).status).toBe(404)
  580. })
  581. it('400s a non-JSON body', async () => {
  582. const response = await handler.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body: 'not json' }))
  583. expect(response.status).toBe(400)
  584. })
  585. it('rejects a malformed envelope with bad-request and the invalid-request sentinel rpcId', async () => {
  586. const response = await handler.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ nope: true }) }))
  587. expect(response.status).toBe(200)
  588. const body = await response.json() as { rpcId: string; result: { ok: boolean; error?: { code: string } } }
  589. expect(body.rpcId).toBe('invalid-request')
  590. expect(body.result.error?.code).toBe('bad-request')
  591. })
  592. it('rejects a method/path mismatch echoing the envelope rpcId', async () => {
  593. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-9', method: 'session.cancel', payload: {} })
  594. const response = await handler.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body }))
  595. const parsed = await response.json() as { rpcId: string; result: { error?: { message: string } } }
  596. expect(parsed.rpcId).toBe('r-9')
  597. expect(parsed.result.error?.message).toContain('does not match path')
  598. })
  599. it('rejects an invalid payload with the zod issues attached', async () => {
  600. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-10', method: 'session.cancel', payload: {} })
  601. const response = await handler.fetch(new Request('http://x/api/session.cancel', { method: 'POST', headers: { 'content-type': 'application/json' }, body }))
  602. const parsed = await response.json() as { result: { error?: { code: string; details: { issues: unknown[] } } } }
  603. expect(parsed.result.error?.code).toBe('bad-request')
  604. expect(parsed.result.error?.details.issues.length).toBeGreaterThan(0)
  605. })
  606. it('500s when the impl itself throws', async () => {
  607. const crashing = toFetchHandler(fakeApi({ crashOn: 'session.list' }))
  608. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-11', method: 'session.list', payload: {} })
  609. const response = await crashing.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body }))
  610. expect(response.status).toBe(500)
  611. expect(await response.text()).toContain('impl crashed')
  612. })
  613. it('routes /api/respond, rejecting malformed client-responses as a receipt', async () => {
  614. const good = JSON.stringify({ type: 'client-response', rpcId: 'known', result: { ok: true, value: null } })
  615. const goodReceipt: unknown = await (await handler.fetch(new Request('http://x/api/respond', { method: 'POST', headers: { 'content-type': 'application/json' }, body: good }))).json()
  616. expect(goodReceipt).toEqual({ accepted: true })
  617. const bad = JSON.stringify({ type: 'client-request', rpcId: 'r', method: 'x', payload: {} })
  618. const badReceipt: unknown = await (await handler.fetch(new Request('http://x/api/respond', { method: 'POST', headers: { 'content-type': 'application/json' }, body: bad }))).json()
  619. expect(badReceipt).toEqual({ accepted: false, reason: 'bad-response' })
  620. })
  621. it('accepts (url, init) form fetch invocation', async () => {
  622. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-12', method: 'session.list', payload: {} })
  623. const response = await handler.fetch('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body })
  624. expect(response.status).toBe(200)
  625. })
  626. })
  627. describe('SSE streams through the carrier', () => {
  628. it('yields mux frames as ServerRequest narrow forms and completes', async () => {
  629. const ac = new AbortController()
  630. const frames = await collect(client().events.mux({}, ac.signal))
  631. expect(frames).toHaveLength(1)
  632. expect(frames[0]?.payload).toMatchObject({ type: 'session/subscribed' })
  633. expect(frames[0]?.rpcId).toBe('frame-0')
  634. })
  635. it('yields host frames', async () => {
  636. const ac = new AbortController()
  637. const frames = await collect(client().events.host({}, ac.signal))
  638. expect(frames[0]?.payload).toMatchObject({ type: 'host/session-removed' })
  639. })
  640. it('drops frames after the consumer aborts mid-stream', async () => {
  641. const many = Array.from({ length: 50 }, (_, i): MuxFrame => ({ type: 'session/subscribed', sessionId: `s${String(i)}` as never, lastSeq: i }))
  642. const ac = new AbortController()
  643. const received: RpcRequest<MuxFrame>[] = []
  644. for await (const envelope of client(fakeApi({ muxFrames: many })).events.mux({}, ac.signal)) {
  645. received.push(envelope)
  646. if (received.length === 2) break // generator return → reader.cancel path
  647. }
  648. expect(received).toHaveLength(2)
  649. })
  650. it('swallows a reader.cancel rejection on early exit', async () => {
  651. const encoder = new TextEncoder()
  652. const body = new ReadableStream<Uint8Array>({
  653. start(controller) {
  654. const frame = { type: 'server-request', rpcId: 'f0', method: 'session/subscribed', payload: { type: 'session/subscribed', sessionId: 's', lastSeq: -1 } }
  655. controller.enqueue(encoder.encode(`data: ${JSON.stringify(frame)}\n\n`))
  656. // stream intentionally left open: the consumer breaks first
  657. },
  658. cancel() {
  659. throw new Error('cancel refused')
  660. },
  661. })
  662. const c = new InProcessApiClient({ fetch: async () => new Response(body, { headers: { 'content-type': 'text/event-stream' } }) })
  663. const received: RpcRequest<MuxFrame>[] = []
  664. for await (const envelope of c.events.mux({}, new AbortController().signal)) {
  665. received.push(envelope)
  666. break
  667. }
  668. expect(received).toHaveLength(1)
  669. })
  670. it('surfaces a mid-stream impl failure as one stream/error frame, then the stream ends', async () => {
  671. const api = fakeApi()
  672. api.events.mux = (_request, _signal) => (async function * (): AsyncGenerator<RpcRequest<MuxFrame>> {
  673. yield { rpcId: RpcId('f0'), payload: { type: 'session/subscribed', sessionId: 's' as never, lastSeq: -1 } }
  674. throw new Error('stream source died')
  675. })()
  676. const frames = await collect(client(api).events.mux({}, new AbortController().signal))
  677. expect(frames).toHaveLength(2)
  678. expect(frames[1]?.payload).toMatchObject({ type: 'stream/error', error: { code: 'internal' } })
  679. })
  680. })
  681. describe('client respond and transport failures', () => {
  682. it('passes a client-response through and parses the receipt', async () => {
  683. const receipt = await client().respond({ type: 'client-response', rpcId: RpcId('known'), result: { ok: true, value: null } })
  684. expect(receipt).toEqual({ accepted: true })
  685. const late = await client().respond({ type: 'client-response', rpcId: RpcId('late'), result: { ok: true, value: null } })
  686. expect(late).toEqual({ accepted: false, reason: 'not-pending' })
  687. })
  688. it('throws on non-OK unary and respond and stream transport', async () => {
  689. const broken = new InProcessApiClient({ fetch: async () => new Response('down', { status: 503 }) })
  690. await expect(broken.sessions.list({})).rejects.toThrow('transport failure for /api/session.list: HTTP 503')
  691. await expect(broken.respond({ type: 'client-response', rpcId: RpcId('r'), result: { ok: true, value: null } }))
  692. .rejects.toThrow('transport failure for /api/respond')
  693. await expect(collect(broken.events.mux({}, new AbortController().signal))).rejects.toThrow('transport failure for /api/events.mux')
  694. })
  695. it('throws on an rpcId echo mismatch', async () => {
  696. const lying = new InProcessApiClient({
  697. fetch: async () => Response.json({ type: 'server-response', rpcId: 'someone-else', result: { ok: true, value: { items: [] } } }),
  698. })
  699. await expect(lying.sessions.list({})).rejects.toThrow('rpcId mismatch')
  700. })
  701. })
  702. describe('envelope observation', () => {
  703. it('batches envelopes per microtask and isolates a throwing listener', async () => {
  704. const c = client()
  705. const batches: (readonly RpcMessage[])[] = []
  706. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  707. const unsubscribeThrowing = c.subscribeEnvelopes(() => { throw new Error('observer bug') })
  708. const unsubscribe = c.subscribeEnvelopes((batch) => { batches.push(batch) })
  709. await c.sessions.list({})
  710. await new Promise((resolve) => { setTimeout(resolve, 0) })
  711. // request and response tap in separate microtask windows (the await between
  712. // them yields), so both arrive but batch count is timing-defined
  713. expect(batches.flatMap(batch => batch.map(message => message.type))).toEqual(['client-request', 'server-response'])
  714. expect(errorSpy).toHaveBeenCalled()
  715. unsubscribe()
  716. unsubscribeThrowing()
  717. errorSpy.mockRestore()
  718. })
  719. it('skips buffering entirely with no listeners and after unsubscribe', async () => {
  720. const c = client()
  721. const seen: RpcMessage[] = []
  722. const unsubscribe = c.subscribeEnvelopes((batch) => { seen.push(...batch) })
  723. unsubscribe()
  724. await c.sessions.list({})
  725. await new Promise((resolve) => { setTimeout(resolve, 0) })
  726. expect(seen).toHaveLength(0)
  727. })
  728. it('coalesces multiple calls in one microtask window into one flush', async () => {
  729. const c = client()
  730. const batches: (readonly RpcMessage[])[] = []
  731. c.subscribeEnvelopes((batch) => { batches.push(batch) })
  732. await Promise.all([c.sessions.list({}), c.host.describe({})])
  733. await new Promise((resolve) => { setTimeout(resolve, 0) })
  734. const total = batches.reduce((n, batch) => n + batch.length, 0)
  735. expect(total).toBe(4)
  736. })
  737. })
  738. describe('resolveBase', () => {
  739. it('prefers a real location.origin and falls back to the internal authority', async () => {
  740. class Probe extends AbstractApiClient {
  741. urls: string[] = []
  742. protected async doFetch(input: URL): Promise<Response> {
  743. this.urls.push(input.href)
  744. return Response.json({ type: 'server-response', rpcId: this.lastMinted, result: { ok: true, value: { items: [] } } })
  745. }
  746. lastMinted = ''
  747. protected override mintRpcId(): ReturnType<AbstractApiClient['mintRpcId']> {
  748. const id = super.mintRpcId()
  749. this.lastMinted = id
  750. return id
  751. }
  752. }
  753. const probe = new Probe()
  754. await probe.sessions.list({})
  755. expect(probe.urls[0]).toMatch(/^http:\/\/dsh\.internal\//)
  756. const globalWithLocation = globalThis as { location?: { origin?: string } }
  757. globalWithLocation.location = { origin: 'http://host.example' }
  758. try {
  759. const probe2 = new Probe()
  760. await probe2.sessions.list({})
  761. expect(probe2.urls[0]).toMatch(/^http:\/\/host\.example\//)
  762. globalWithLocation.location = { origin: 'null' } // sandboxed iframe shape
  763. const probe3 = new Probe()
  764. await probe3.sessions.list({})
  765. expect(probe3.urls[0]).toMatch(/^http:\/\/dsh\.internal\//)
  766. } finally {
  767. delete globalWithLocation.location
  768. }
  769. })
  770. })