fetch-carrier.spec.ts 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770
  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. skills: {
  196. async list(request) {
  197. return { rpcId: request.rpcId, result: { ok: true, value: { skills: [{ name: 'commit-helper', description: 'Git commits' }] } } }
  198. },
  199. },
  200. goals: {
  201. async create(request) {
  202. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  203. },
  204. async edit(request) {
  205. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  206. },
  207. async pause(request) {
  208. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  209. },
  210. async resume(request) {
  211. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  212. },
  213. async complete(request) {
  214. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  215. },
  216. async clear(request) {
  217. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'internal', message: 'stub', details: {} } } }
  218. },
  219. },
  220. settings: {
  221. async describe(request) {
  222. return { rpcId: request.rpcId, result: { ok: true, value: { writable: true, hasDocument: false, namespaces: [] } } }
  223. },
  224. async openDocument(request) {
  225. return { rpcId: request.rpcId, result: { ok: true, value: { opened: true as const } } }
  226. },
  227. async update(request) {
  228. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'settings-rejected', message: 'stub', details: { ns: request.payload.ns } } } }
  229. },
  230. async replace(request) {
  231. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'settings-rejected', message: 'stub', details: { ns: request.payload.ns } } } }
  232. },
  233. async mutate(request) {
  234. return { rpcId: request.rpcId, result: { ok: false, error: { code: 'settings-rejected', message: 'stub', details: { ns: request.payload.ns } } } }
  235. },
  236. },
  237. credentials: {
  238. async describe(request) {
  239. return { rpcId: request.rpcId, result: { ok: true, value: { credentials: {} } } }
  240. },
  241. async set(request) {
  242. return { rpcId: request.rpcId, result: { ok: true, value: {} } }
  243. },
  244. async unset(request) {
  245. return { rpcId: request.rpcId, result: { ok: true, value: {} } }
  246. },
  247. },
  248. llm: {
  249. async providers(request) {
  250. return { rpcId: request.rpcId, result: { ok: true, value: { providers: [] } } }
  251. },
  252. async models(request) {
  253. return { rpcId: request.rpcId, result: { ok: true, value: { groups: [], failures: [] } } }
  254. },
  255. },
  256. events: {
  257. mux: (_request, signal) => stream(muxFrames, signal),
  258. host: (_request, signal) => stream(hostFrames, signal),
  259. },
  260. async respond(message: ClientResponse): Promise<RpcReceipt> {
  261. return message.rpcId === 'known' ? { accepted: true } : { accepted: false, reason: 'not-pending' }
  262. },
  263. }
  264. }
  265. function client(api: ApiProxy = fakeApi(), timeoutMs?: number): InProcessApiClient {
  266. return new InProcessApiClient(toFetchHandler(api), timeoutMs)
  267. }
  268. async function collect<F>(stream: AsyncIterable<RpcRequest<F>>): Promise<RpcRequest<F>[]> {
  269. const out: RpcRequest<F>[] = []
  270. for await (const envelope of stream) out.push(envelope)
  271. return out
  272. }
  273. describe('unary round trip (handler ⇄ client, no network)', () => {
  274. it('carries a success result and echoes the minted rpcId', async () => {
  275. const response = await client().sessions.list({})
  276. expect(response.result).toEqual({ ok: true, value: { items: [] } })
  277. expect(response.rpcId).toMatch(/[0-9a-f-]{36}/)
  278. })
  279. it('carries the tail-page projections block through the wire schema (Zod must not strip it)', async () => {
  280. const response = await client().sessions.history({ sessionId: 'with-projections' as never })
  281. expect(response.result.ok).toBe(true)
  282. if (response.result.ok) {
  283. expect(response.result.value.projections).toEqual(
  284. { asOfSeq: 9, values: { todos: [{ content: 'current', status: 'in_progress' }] } },
  285. )
  286. }
  287. })
  288. it('carries a business error as 200 + error result', async () => {
  289. const response = await client().sessions.history({ sessionId: 'missing' as never })
  290. expect(response.result.ok).toBe(false)
  291. if (!response.result.ok) expect(response.result.error.code).toBe('session-not-found')
  292. })
  293. it('covers create/prompt/updateQueue/cancel/describe passthrough', async () => {
  294. const c = client()
  295. expect((await c.sessions.search({ query: 'fixture' })).result).toEqual({
  296. ok: true,
  297. value: { items: [{ sessionId: 's1', snippet: 'fixture match' }], hasMore: false },
  298. })
  299. expect((await c.sessions.create({})).result.ok).toBe(true)
  300. expect((await c.sessions.models({ sessionId: 's' as never })).result.ok).toBe(true)
  301. const selected = await c.sessions.selectModel({
  302. sessionId: 's' as never,
  303. provider: 'deepseek-official',
  304. model: 'deepseek-v4-flash',
  305. reasoningEffort: 'max',
  306. })
  307. expect(selected.result).toMatchObject({
  308. ok: true,
  309. value: {
  310. selected: {
  311. provider: 'deepseek-official',
  312. model: 'deepseek-v4-flash',
  313. reasoningEffort: 'max',
  314. },
  315. },
  316. })
  317. const renamed = await c.sessions.rename({ sessionId: 's' as never, title: 'named' })
  318. expect(renamed.result).toMatchObject({ ok: true, value: { title: 'named', seq: 0 } })
  319. expect((await c.sessions.prompt({ sessionId: 's' as never, mode: 'queue', content: [{ type: 'text', text: 'x' }] })).result.ok).toBe(true)
  320. expect((await c.sessions.updateQueue({
  321. sessionId: 's' as never,
  322. itemId: 'item-1' as never,
  323. action: { kind: 'remove' },
  324. })).result.ok).toBe(true)
  325. expect((await c.sessions.cancel({ sessionId: 's' as never })).result.ok).toBe(true)
  326. expect((await c.host.describe({})).result.ok).toBe(true)
  327. })
  328. it('round-trips the native picker without the default unary timeout', async () => {
  329. const api = fakeApi()
  330. api.host.pickDirectory = async (request) => {
  331. await new Promise(resolve => setTimeout(resolve, 15))
  332. return { rpcId: request.rpcId, result: { ok: true, value: { path: '/tmp/project' } } }
  333. }
  334. const response = await client(api, 1).host.pickDirectory({})
  335. expect(response.result).toEqual({ ok: true, value: { path: '/tmp/project' } })
  336. })
  337. it('round-trips the browse listing and creation calls through the wire form', async () => {
  338. const c = client()
  339. const listed = await c.host.listDirectory({ path: '/w' })
  340. expect(listed.result).toEqual({
  341. ok: true,
  342. value: { path: '/w', home: '/w', crumbs: [{ name: '/', path: '/', hidden: false }], entries: [], truncated: false },
  343. })
  344. const home = await c.host.listDirectory({})
  345. expect(home.result).toMatchObject({ ok: true, value: { home: '/w' } })
  346. const created = await c.host.createDirectory({ path: '/w', name: 'fresh' })
  347. expect(created.result).toEqual({ ok: true, value: { path: '/w/new' } })
  348. })
  349. it('round-trips host.openPath through the wire form', async () => {
  350. const api = fakeApi()
  351. let opened: string | undefined
  352. api.host.openPath = async (request) => {
  353. opened = request.payload.path
  354. return { rpcId: request.rpcId, result: { ok: true, value: { opened: true as const } } }
  355. }
  356. const response = await client(api).host.openPath({ path: '/tmp/a.txt' })
  357. expect(opened).toBe('/tmp/a.txt')
  358. expect(response.result).toEqual({ ok: true, value: { opened: true } })
  359. })
  360. it('round-trips command.list / command.execute / skill.list through the wire form', async () => {
  361. const c = client()
  362. const list = await c.commands.list({ sessionId: 's' as never })
  363. expect(list.result).toEqual({ ok: true, value: { commands: [{ name: 'plan', description: 'Toggle plan mode', input: { hint: 'on|off' } }] } })
  364. const hit = await c.commands.execute({ sessionId: 's' as never, line: '/plan off' })
  365. expect(hit.result).toEqual({ ok: true, value: { matched: true, commandId: 'cmd-x' } })
  366. const miss = await c.commands.execute({ sessionId: 's' as never, line: '/nope' })
  367. expect(miss.result).toEqual({ ok: true, value: { matched: false } })
  368. const skills = await c.skills.list({ sessionId: 's' as never })
  369. expect(skills.result).toEqual({ ok: true, value: { skills: [{ name: 'commit-helper', description: 'Git commits' }] } })
  370. })
  371. it('lets command.execute finish after the 30-second default unary deadline', async () => {
  372. vi.useFakeTimers()
  373. const timeoutSpy = vi.spyOn(AbortSignal, 'timeout').mockImplementation((milliseconds) => {
  374. const controller = new AbortController()
  375. setTimeout(() => {
  376. controller.abort(new DOMException('The operation was aborted due to timeout', 'TimeoutError'))
  377. }, milliseconds)
  378. return controller.signal
  379. })
  380. try {
  381. const api = fakeApi()
  382. api.commands.execute = async (request) => {
  383. await new Promise(resolve => setTimeout(resolve, 30_001))
  384. return {
  385. rpcId: request.rpcId,
  386. result: { ok: true, value: { matched: true, commandId: CommandId('cmd-slow') } },
  387. }
  388. }
  389. const execution = client(api).commands.execute({ sessionId: 's' as never, line: '/slow' })
  390. const assertion = expect(execution).resolves.toMatchObject({
  391. result: { ok: true, value: { matched: true, commandId: 'cmd-slow' } },
  392. })
  393. await Promise.all([
  394. vi.advanceTimersByTimeAsync(30_001),
  395. assertion,
  396. ])
  397. expect(timeoutSpy).not.toHaveBeenCalled()
  398. } finally {
  399. timeoutSpy.mockRestore()
  400. vi.useRealTimers()
  401. }
  402. })
  403. it('round-trips the subagent domain through the wire form', async () => {
  404. const c = client()
  405. expect((await c.subagents.list({ parentSessionId: 'parent' as never })).result)
  406. .toEqual({ ok: true, value: { entries: [], parentAvailable: false } })
  407. expect((await c.subagents.history({
  408. parentSessionId: 'parent' as never,
  409. childSessionId: 'child' as never,
  410. mode: 'one-shot',
  411. })).result).toEqual({ ok: true, value: { events: [], hasMore: false } })
  412. expect((await c.subagents.prompt({
  413. parentSessionId: 'parent' as never,
  414. childSessionId: 'child' as never,
  415. mode: 'continuable',
  416. content: [],
  417. })).result).toEqual({ ok: true, value: { messageId: 'message-1' } })
  418. })
  419. it('keeps caller and connection aborts on command.execute', async () => {
  420. const api = fakeApi()
  421. const started = Promise.withResolvers<AbortSignal>()
  422. api.commands.execute = async (request, signal) => {
  423. started.resolve(signal)
  424. if (!signal.aborted) {
  425. await new Promise<void>((resolve) => {
  426. signal.addEventListener('abort', () => { resolve() }, { once: true })
  427. })
  428. }
  429. return {
  430. rpcId: request.rpcId,
  431. result: { ok: false, error: { code: 'cancelled', message: 'aborted', details: {} } },
  432. }
  433. }
  434. const controller = new AbortController()
  435. const execution = client(api).commands.execute(
  436. { sessionId: 's' as never, line: '/hang' },
  437. controller.signal,
  438. )
  439. const handlerSignal = await started.promise
  440. controller.abort(new Error('connection closed'))
  441. await expect(execution).rejects.toThrow('connection closed')
  442. expect(handlerSignal.aborted).toBe(true)
  443. })
  444. it('propagates the carrier Request signal into command.execute', async () => {
  445. const handler = toFetchHandler(fakeApi())
  446. const controller = new AbortController()
  447. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-sig', method: 'command.execute', payload: { sessionId: 's', line: '/hang' } })
  448. // The fake's /hang settles only when the invoke-level signal aborts: a
  449. // completed response with the cancelled error proves req.signal reached it.
  450. const pending = handler.fetch(new Request('http://x/api/command.execute', { method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal }))
  451. controller.abort()
  452. const response = await pending
  453. const parsed = await response.json() as { rpcId: string; result: { ok: boolean; error?: { code: string } } }
  454. expect(parsed.rpcId).toBe('r-sig')
  455. expect(parsed.result.error?.code).toBe('cancelled')
  456. })
  457. it('propagates the carrier Request signal into session.search', async () => {
  458. const handler = toFetchHandler(fakeApi())
  459. const controller = new AbortController()
  460. const body = JSON.stringify({
  461. type: 'client-request',
  462. rpcId: 'r-search-sig',
  463. method: 'session.search',
  464. payload: { query: 'hang' },
  465. })
  466. const pending = handler.fetch(new Request(
  467. 'http://x/api/session.search',
  468. { method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal },
  469. ))
  470. controller.abort()
  471. const response = await pending
  472. const parsed = await response.json() as {
  473. rpcId: string
  474. result: { error?: { code: string } }
  475. }
  476. expect(parsed.rpcId).toBe('r-search-sig')
  477. expect(parsed.result.error?.code).toBe('cancelled')
  478. })
  479. it('propagates the carrier Request signal into subagent.prompt', async () => {
  480. const handler = toFetchHandler(fakeApi())
  481. const controller = new AbortController()
  482. const body = JSON.stringify({
  483. type: 'client-request',
  484. rpcId: 'r-subagent-sig',
  485. method: 'subagent.prompt',
  486. payload: {
  487. parentSessionId: 'parent',
  488. childSessionId: 'child',
  489. mode: 'continuable',
  490. content: [{ type: 'text', text: 'hang' }],
  491. },
  492. })
  493. const pending = handler.fetch(new Request(
  494. 'http://x/api/subagent.prompt',
  495. { method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal },
  496. ))
  497. controller.abort()
  498. const response = await pending
  499. const parsed = await response.json() as {
  500. rpcId: string
  501. result: { error?: { code: string } }
  502. }
  503. expect(parsed.rpcId).toBe('r-subagent-sig')
  504. expect(parsed.result.error?.code).toBe('cancelled')
  505. })
  506. it('propagates the carrier Request signal into host.pickDirectory', async () => {
  507. const api = fakeApi()
  508. api.host.pickDirectory = async (request, signal) => {
  509. if (!signal.aborted) {
  510. await new Promise<void>((resolve) => {
  511. signal.addEventListener('abort', () => { resolve() }, { once: true })
  512. })
  513. }
  514. return {
  515. rpcId: request.rpcId,
  516. result: { ok: false, error: { code: 'cancelled', message: 'aborted', details: {} } },
  517. }
  518. }
  519. const handler = toFetchHandler(api)
  520. const controller = new AbortController()
  521. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-picker', method: 'host.pickDirectory', payload: {} })
  522. const pending = handler.fetch(new Request('http://x/api/host.pickDirectory', {
  523. method: 'POST', headers: { 'content-type': 'application/json' }, body, signal: controller.signal,
  524. }))
  525. controller.abort()
  526. const parsed = await (await pending).json() as { result: { error?: { code: string } } }
  527. expect(parsed.result.error?.code).toBe('cancelled')
  528. })
  529. })
  530. describe('handler carrier-layer statuses', () => {
  531. const handler = toFetchHandler(fakeApi())
  532. it('404s unknown paths and non-POST non-stream methods', async () => {
  533. expect((await handler.fetch(new Request('http://x/other', { method: 'POST', headers: { 'content-type': 'application/json' }, body: '{}' }))).status).toBe(404)
  534. expect((await handler.fetch(new Request('http://x/api/session.list', { method: 'GET' }))).status).toBe(404)
  535. 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)
  536. })
  537. it('400s a non-JSON body', async () => {
  538. const response = await handler.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body: 'not json' }))
  539. expect(response.status).toBe(400)
  540. })
  541. it('rejects a malformed envelope with bad-request and the invalid-request sentinel rpcId', async () => {
  542. const response = await handler.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ nope: true }) }))
  543. expect(response.status).toBe(200)
  544. const body = await response.json() as { rpcId: string; result: { ok: boolean; error?: { code: string } } }
  545. expect(body.rpcId).toBe('invalid-request')
  546. expect(body.result.error?.code).toBe('bad-request')
  547. })
  548. it('rejects a method/path mismatch echoing the envelope rpcId', async () => {
  549. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-9', method: 'session.cancel', payload: {} })
  550. const response = await handler.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body }))
  551. const parsed = await response.json() as { rpcId: string; result: { error?: { message: string } } }
  552. expect(parsed.rpcId).toBe('r-9')
  553. expect(parsed.result.error?.message).toContain('does not match path')
  554. })
  555. it('rejects an invalid payload with the zod issues attached', async () => {
  556. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-10', method: 'session.cancel', payload: {} })
  557. const response = await handler.fetch(new Request('http://x/api/session.cancel', { method: 'POST', headers: { 'content-type': 'application/json' }, body }))
  558. const parsed = await response.json() as { result: { error?: { code: string; details: { issues: unknown[] } } } }
  559. expect(parsed.result.error?.code).toBe('bad-request')
  560. expect(parsed.result.error?.details.issues.length).toBeGreaterThan(0)
  561. })
  562. it('500s when the impl itself throws', async () => {
  563. const crashing = toFetchHandler(fakeApi({ crashOn: 'session.list' }))
  564. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-11', method: 'session.list', payload: {} })
  565. const response = await crashing.fetch(new Request('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body }))
  566. expect(response.status).toBe(500)
  567. expect(await response.text()).toContain('impl crashed')
  568. })
  569. it('routes /api/respond, rejecting malformed client-responses as a receipt', async () => {
  570. const good = JSON.stringify({ type: 'client-response', rpcId: 'known', result: { ok: true, value: null } })
  571. const goodReceipt: unknown = await (await handler.fetch(new Request('http://x/api/respond', { method: 'POST', headers: { 'content-type': 'application/json' }, body: good }))).json()
  572. expect(goodReceipt).toEqual({ accepted: true })
  573. const bad = JSON.stringify({ type: 'client-request', rpcId: 'r', method: 'x', payload: {} })
  574. const badReceipt: unknown = await (await handler.fetch(new Request('http://x/api/respond', { method: 'POST', headers: { 'content-type': 'application/json' }, body: bad }))).json()
  575. expect(badReceipt).toEqual({ accepted: false, reason: 'bad-response' })
  576. })
  577. it('accepts (url, init) form fetch invocation', async () => {
  578. const body = JSON.stringify({ type: 'client-request', rpcId: 'r-12', method: 'session.list', payload: {} })
  579. const response = await handler.fetch('http://x/api/session.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body })
  580. expect(response.status).toBe(200)
  581. })
  582. })
  583. describe('SSE streams through the carrier', () => {
  584. it('yields mux frames as ServerRequest narrow forms and completes', async () => {
  585. const ac = new AbortController()
  586. const frames = await collect(client().events.mux({}, ac.signal))
  587. expect(frames).toHaveLength(1)
  588. expect(frames[0]?.payload).toMatchObject({ type: 'session/subscribed' })
  589. expect(frames[0]?.rpcId).toBe('frame-0')
  590. })
  591. it('yields host frames', async () => {
  592. const ac = new AbortController()
  593. const frames = await collect(client().events.host({}, ac.signal))
  594. expect(frames[0]?.payload).toMatchObject({ type: 'host/session-removed' })
  595. })
  596. it('drops frames after the consumer aborts mid-stream', async () => {
  597. const many = Array.from({ length: 50 }, (_, i): MuxFrame => ({ type: 'session/subscribed', sessionId: `s${String(i)}` as never, lastSeq: i }))
  598. const ac = new AbortController()
  599. const received: RpcRequest<MuxFrame>[] = []
  600. for await (const envelope of client(fakeApi({ muxFrames: many })).events.mux({}, ac.signal)) {
  601. received.push(envelope)
  602. if (received.length === 2) break // generator return → reader.cancel path
  603. }
  604. expect(received).toHaveLength(2)
  605. })
  606. it('swallows a reader.cancel rejection on early exit', async () => {
  607. const encoder = new TextEncoder()
  608. const body = new ReadableStream<Uint8Array>({
  609. start(controller) {
  610. const frame = { type: 'server-request', rpcId: 'f0', method: 'session/subscribed', payload: { type: 'session/subscribed', sessionId: 's', lastSeq: -1 } }
  611. controller.enqueue(encoder.encode(`data: ${JSON.stringify(frame)}\n\n`))
  612. // stream intentionally left open: the consumer breaks first
  613. },
  614. cancel() {
  615. throw new Error('cancel refused')
  616. },
  617. })
  618. const c = new InProcessApiClient({ fetch: async () => new Response(body, { headers: { 'content-type': 'text/event-stream' } }) })
  619. const received: RpcRequest<MuxFrame>[] = []
  620. for await (const envelope of c.events.mux({}, new AbortController().signal)) {
  621. received.push(envelope)
  622. break
  623. }
  624. expect(received).toHaveLength(1)
  625. })
  626. it('surfaces a mid-stream impl failure as one stream/error frame, then the stream ends', async () => {
  627. const api = fakeApi()
  628. api.events.mux = (_request, _signal) => (async function * (): AsyncGenerator<RpcRequest<MuxFrame>> {
  629. yield { rpcId: RpcId('f0'), payload: { type: 'session/subscribed', sessionId: 's' as never, lastSeq: -1 } }
  630. throw new Error('stream source died')
  631. })()
  632. const frames = await collect(client(api).events.mux({}, new AbortController().signal))
  633. expect(frames).toHaveLength(2)
  634. expect(frames[1]?.payload).toMatchObject({ type: 'stream/error', error: { code: 'internal' } })
  635. })
  636. })
  637. describe('client respond and transport failures', () => {
  638. it('passes a client-response through and parses the receipt', async () => {
  639. const receipt = await client().respond({ type: 'client-response', rpcId: RpcId('known'), result: { ok: true, value: null } })
  640. expect(receipt).toEqual({ accepted: true })
  641. const late = await client().respond({ type: 'client-response', rpcId: RpcId('late'), result: { ok: true, value: null } })
  642. expect(late).toEqual({ accepted: false, reason: 'not-pending' })
  643. })
  644. it('throws on non-OK unary and respond and stream transport', async () => {
  645. const broken = new InProcessApiClient({ fetch: async () => new Response('down', { status: 503 }) })
  646. await expect(broken.sessions.list({})).rejects.toThrow('transport failure for /api/session.list: HTTP 503')
  647. await expect(broken.respond({ type: 'client-response', rpcId: RpcId('r'), result: { ok: true, value: null } }))
  648. .rejects.toThrow('transport failure for /api/respond')
  649. await expect(collect(broken.events.mux({}, new AbortController().signal))).rejects.toThrow('transport failure for /api/events.mux')
  650. })
  651. it('throws on an rpcId echo mismatch', async () => {
  652. const lying = new InProcessApiClient({
  653. fetch: async () => Response.json({ type: 'server-response', rpcId: 'someone-else', result: { ok: true, value: { items: [] } } }),
  654. })
  655. await expect(lying.sessions.list({})).rejects.toThrow('rpcId mismatch')
  656. })
  657. })
  658. describe('envelope observation', () => {
  659. it('batches envelopes per microtask and isolates a throwing listener', async () => {
  660. const c = client()
  661. const batches: (readonly RpcMessage[])[] = []
  662. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  663. const unsubscribeThrowing = c.subscribeEnvelopes(() => { throw new Error('observer bug') })
  664. const unsubscribe = c.subscribeEnvelopes((batch) => { batches.push(batch) })
  665. await c.sessions.list({})
  666. await new Promise((resolve) => { setTimeout(resolve, 0) })
  667. // request and response tap in separate microtask windows (the await between
  668. // them yields), so both arrive but batch count is timing-defined
  669. expect(batches.flatMap(batch => batch.map(message => message.type))).toEqual(['client-request', 'server-response'])
  670. expect(errorSpy).toHaveBeenCalled()
  671. unsubscribe()
  672. unsubscribeThrowing()
  673. errorSpy.mockRestore()
  674. })
  675. it('skips buffering entirely with no listeners and after unsubscribe', async () => {
  676. const c = client()
  677. const seen: RpcMessage[] = []
  678. const unsubscribe = c.subscribeEnvelopes((batch) => { seen.push(...batch) })
  679. unsubscribe()
  680. await c.sessions.list({})
  681. await new Promise((resolve) => { setTimeout(resolve, 0) })
  682. expect(seen).toHaveLength(0)
  683. })
  684. it('coalesces multiple calls in one microtask window into one flush', async () => {
  685. const c = client()
  686. const batches: (readonly RpcMessage[])[] = []
  687. c.subscribeEnvelopes((batch) => { batches.push(batch) })
  688. await Promise.all([c.sessions.list({}), c.host.describe({})])
  689. await new Promise((resolve) => { setTimeout(resolve, 0) })
  690. const total = batches.reduce((n, batch) => n + batch.length, 0)
  691. expect(total).toBe(4)
  692. })
  693. })
  694. describe('resolveBase', () => {
  695. it('prefers a real location.origin and falls back to the internal authority', async () => {
  696. class Probe extends AbstractApiClient {
  697. urls: string[] = []
  698. protected async doFetch(input: URL): Promise<Response> {
  699. this.urls.push(input.href)
  700. return Response.json({ type: 'server-response', rpcId: this.lastMinted, result: { ok: true, value: { items: [] } } })
  701. }
  702. lastMinted = ''
  703. protected override mintRpcId(): ReturnType<AbstractApiClient['mintRpcId']> {
  704. const id = super.mintRpcId()
  705. this.lastMinted = id
  706. return id
  707. }
  708. }
  709. const probe = new Probe()
  710. await probe.sessions.list({})
  711. expect(probe.urls[0]).toMatch(/^http:\/\/dsh\.internal\//)
  712. const globalWithLocation = globalThis as { location?: { origin?: string } }
  713. globalWithLocation.location = { origin: 'http://host.example' }
  714. try {
  715. const probe2 = new Probe()
  716. await probe2.sessions.list({})
  717. expect(probe2.urls[0]).toMatch(/^http:\/\/host\.example\//)
  718. globalWithLocation.location = { origin: 'null' } // sandboxed iframe shape
  719. const probe3 = new Probe()
  720. await probe3.sessions.list({})
  721. expect(probe3.urls[0]).toMatch(/^http:\/\/dsh\.internal\//)
  722. } finally {
  723. delete globalWithLocation.location
  724. }
  725. })
  726. })