client-handler.spec.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407
  1. /**
  2. * Wire-protocol coverage over the isomorphic point: InProcessApiClient →
  3. * toFetchHandler(scripted impl) runs the real envelope wrap/unwrap, zod
  4. * two-level parse, rpcId discipline, and SSE framing with no network and no
  5. * browser. Each case scripts its own minimal ApiProxy.
  6. */
  7. import { describe, expect, it, vi } from 'vitest'
  8. import type { SessionId } from '@deepseek-ai/dsh-session'
  9. import type { ApiProxy, HostFrame, MuxFrame, RpcMessage, RpcRequest, RpcResponse } from '@deepseek-ai/dsh-host-apiproxy'
  10. import { InProcessApiClient, RpcId, toFetchHandler } from '@deepseek-ai/dsh-host-apiproxy'
  11. const sid = (id: string): SessionId => id as SessionId
  12. function ok<T>(request: RpcRequest<unknown>, value: T): Promise<RpcResponse<T>> {
  13. return Promise.resolve({ rpcId: request.rpcId, result: { ok: true, value } })
  14. }
  15. /** Scripted impl: every method resolves an empty-ish OK unless a case overrides it. */
  16. function scriptedApi(overrides: {
  17. sessions?: Partial<ApiProxy['sessions']>
  18. host?: Partial<ApiProxy['host']>
  19. events?: Partial<ApiProxy['events']>
  20. respond?: ApiProxy['respond']
  21. } = {}): ApiProxy {
  22. async function *empty<F>(): AsyncGenerator<RpcRequest<F>> { /* no frames */ }
  23. return {
  24. sessions: {
  25. list: r => ok(r, { items: [] }),
  26. create: r => ok(r, { sessionId: sid('s-new') }),
  27. history: r => ok(r, { events: [], hasMore: false }),
  28. prompt: r => ok(r, { accepted: true as const }),
  29. cancel: r => ok(r, { accepted: true as const }),
  30. ...overrides.sessions,
  31. },
  32. host: { describe: r => ok(r, { version: '0-test', cwd: '/t', attachedSessions: 0 }), ...overrides.host },
  33. events: { mux: () => empty<MuxFrame>(), host: () => empty<HostFrame>(), ...overrides.events },
  34. respond: overrides.respond ?? (() => Promise.resolve({ accepted: false as const, reason: 'not-pending' as const })),
  35. }
  36. }
  37. function client(api: ApiProxy, timeoutMs?: number): InProcessApiClient {
  38. return new InProcessApiClient(toFetchHandler(api), timeoutMs)
  39. }
  40. describe('unary round trip', () => {
  41. it('carries payload out and value back through the full wire form', async () => {
  42. let seen: RpcRequest<{ cursor?: string }> | undefined
  43. const api = scriptedApi({
  44. sessions: {
  45. list: (r) => {
  46. seen = r
  47. return ok(r, { items: [{ sessionId: sid('s1'), updatedAt: 7, running: false }] })
  48. },
  49. },
  50. })
  51. const response = await client(api).sessions.list({ cursor: 'c1' })
  52. // Impl received the narrow form with a minted id; client returned the same id and value.
  53. expect(seen?.payload).toEqual({ cursor: 'c1' })
  54. expect(seen?.rpcId).toBeTruthy()
  55. expect(response.rpcId).toBe(seen?.rpcId)
  56. expect(response.result).toEqual({ ok: true, value: { items: [{ sessionId: 's1', updatedAt: 7, running: false }] } })
  57. })
  58. it('passes business errors through as 200 + err result, not a throw', async () => {
  59. const api = scriptedApi({
  60. sessions: {
  61. cancel: r => Promise.resolve({ rpcId: r.rpcId, result: { ok: false, error: { code: 'session-not-found', message: 'nope', details: { sessionId: sid('sx') } } } }),
  62. },
  63. })
  64. const response = await client(api).sessions.cancel({ sessionId: sid('sx') })
  65. expect(response.result).toEqual({ ok: false, error: { code: 'session-not-found', message: 'nope', details: { sessionId: 'sx' } } })
  66. })
  67. it('throws on rpcId echo mismatch', async () => {
  68. const api = scriptedApi({
  69. sessions: { list: () => Promise.resolve({ rpcId: RpcId('forged'), result: { ok: true, value: { items: [] } } }) },
  70. })
  71. await expect(client(api).sessions.list({})).rejects.toThrow(/rpcId mismatch/)
  72. })
  73. it('rejects an invalid payload at the handler as 200 + bad-request with issues', async () => {
  74. const api = scriptedApi()
  75. const response = await client(api).sessions.history({ sessionId: 123 as unknown as SessionId })
  76. expect(response.result.ok).toBe(false)
  77. if (!response.result.ok) {
  78. expect(response.result.error.code).toBe('bad-request')
  79. expect((response.result.error.details as { issues: unknown[] }).issues.length).toBeGreaterThan(0)
  80. }
  81. })
  82. it('rejects a method/path mismatch as bad-request', async () => {
  83. const handler = toFetchHandler(scriptedApi())
  84. const body = { type: 'client-request', rpcId: 'r1', method: 'session.create', payload: {} }
  85. const response = await handler.fetch('http://dsh.internal/api/session.list', { method: 'POST', body: JSON.stringify(body) })
  86. expect(response.status).toBe(200)
  87. const parsed = await response.json() as { result: { ok: boolean; error?: { code: string; message: string } } }
  88. expect(parsed.result.ok).toBe(false)
  89. expect(parsed.result.error?.code).toBe('bad-request')
  90. expect(parsed.result.error?.message).toMatch(/does not match path/)
  91. })
  92. it('rejects a malformed envelope as bad-request, salvaging the rpcId or falling back to the sentinel', async () => {
  93. const handler = toFetchHandler(scriptedApi())
  94. // No salvageable rpcId → the fixed invalid-request sentinel keeps the response a valid ServerResponse.
  95. const noId = await handler.fetch('http://dsh.internal/api/session.list', { method: 'POST', body: JSON.stringify({ nonsense: true }) })
  96. expect(noId.status).toBe(200)
  97. const noIdParsed = await noId.json() as { rpcId: string; result: { ok: boolean } }
  98. expect(noIdParsed.result.ok).toBe(false)
  99. expect(noIdParsed.rpcId).toBe('invalid-request')
  100. // A string rpcId in the otherwise-bad body is salvaged for correlation.
  101. const withId = await handler.fetch('http://dsh.internal/api/session.list', { method: 'POST', body: JSON.stringify({ rpcId: 'salvage-me', nonsense: true }) })
  102. const withIdParsed = await withId.json() as { rpcId: string; result: { ok: boolean } }
  103. expect(withIdParsed.result.ok).toBe(false)
  104. expect(withIdParsed.rpcId).toBe('salvage-me')
  105. })
  106. it('maps carrier failures to HTTP statuses and the client throws transport failure', async () => {
  107. const handler = toFetchHandler(scriptedApi())
  108. // Unknown method → 404.
  109. const notFound = await handler.fetch('http://dsh.internal/api/no.such', { method: 'POST', body: '{}' })
  110. expect(notFound.status).toBe(404)
  111. // Non-JSON body → 400.
  112. const badBody = await handler.fetch('http://dsh.internal/api/session.list', { method: 'POST', body: '{oops' })
  113. expect(badBody.status).toBe(400)
  114. // Impl crash → 500, and through the client that is a throw, not an err result.
  115. const crashing = scriptedApi({ sessions: { list: () => { throw new Error('impl exploded') } } })
  116. await expect(client(crashing).sessions.list({})).rejects.toThrow(/transport failure .*500/)
  117. })
  118. it('rejects when the transport never resolves within timeoutMs', async () => {
  119. // AbortSignal.timeout is immune to fake timers; a short real timeout keeps this fast.
  120. const never = new InProcessApiClient({
  121. fetch: (_i: RequestInfo | URL, init?: RequestInit) => new Promise<Response>((_resolve, reject) => {
  122. init?.signal?.addEventListener('abort', () => { reject(new Error('aborted by timeout')) })
  123. }),
  124. }, 25)
  125. await expect(never.sessions.list({})).rejects.toThrow()
  126. })
  127. it('aborts a unary call through the caller-supplied external signal', async () => {
  128. // Real-fetch semantics: on abort the rejection is the signal's reason, and the abort
  129. // works even when the transport ignores the signal entirely (hung impl).
  130. const gate = new AbortController()
  131. const hung = new InProcessApiClient({ fetch: () => new Promise<Response>(() => {}) }, 60_000)
  132. const call = hung.sessions.list({}, gate.signal)
  133. gate.abort(new Error('externally aborted'))
  134. await expect(call).rejects.toThrow(/externally aborted/)
  135. })
  136. it('rejects an already-aborted signal before touching the transport, mapping a string reason to an Error', async () => {
  137. let touched = false
  138. const c = new InProcessApiClient({
  139. fetch: () => {
  140. touched = true
  141. return Promise.resolve(new Response('{}'))
  142. },
  143. }, 60_000)
  144. const gate = new AbortController()
  145. gate.abort('gone before start')
  146. await expect(c.sessions.list({}, gate.signal)).rejects.toThrow('gone before start')
  147. expect(touched).toBe(false)
  148. })
  149. it('maps a non-Error, non-string abort reason to the default AbortError message', async () => {
  150. const gate = new AbortController()
  151. const hung = new InProcessApiClient({ fetch: () => new Promise<Response>(() => {}) }, 60_000)
  152. const call = hung.sessions.list({}, gate.signal)
  153. gate.abort(42)
  154. await expect(call).rejects.toThrow('This operation was aborted')
  155. })
  156. it('passes a signal-less doFetch straight through to the handler', async () => {
  157. class Probe extends InProcessApiClient {
  158. direct(url: URL): Promise<Response> {
  159. return this.doFetch(url)
  160. }
  161. }
  162. const probe = new Probe({ fetch: () => Promise.resolve(new Response('raw')) })
  163. const response = await probe.direct(new URL('http://dsh.internal/probe'))
  164. expect(await response.text()).toBe('raw')
  165. })
  166. it('throws on an S→C ok value that fails the method value schema (second-level parse)', async () => {
  167. // Impl echoes rpcId but returns a wrong-shaped value: envelope parse passes, value parse must reject.
  168. const api = scriptedApi({
  169. sessions: { list: r => Promise.resolve({ rpcId: r.rpcId, result: { ok: true, value: { items: 'not-an-array' } } }) as never },
  170. })
  171. await expect(client(api).sessions.list({})).rejects.toThrow()
  172. })
  173. })
  174. describe('SSE stream path', () => {
  175. it('yields frames in order and skips the comment preamble', async () => {
  176. const frames: MuxFrame[] = [
  177. { type: 'session/subscribed', sessionId: sid('s1'), lastSeq: 3 },
  178. { type: 'stream/error', error: { code: 'internal', message: 'x', details: {} } },
  179. ]
  180. const api = scriptedApi({
  181. events: {
  182. async *mux(request) {
  183. let n = 0
  184. for (const frame of frames) yield { rpcId: RpcId(`push-${n++}-${request.rpcId}`), payload: frame }
  185. },
  186. },
  187. })
  188. const seen: MuxFrame[] = []
  189. for await (const envelope of client(api).events.mux({}, new AbortController().signal)) {
  190. seen.push(envelope.payload)
  191. }
  192. expect(seen).toEqual(frames)
  193. })
  194. it('reassembles frames across arbitrary chunk boundaries', async () => {
  195. // Two SSE frames split so one frame spans chunks and one chunk carries parts of both.
  196. const f1 = { type: 'server-request', rpcId: 'a', method: 'session/subscribed', payload: { type: 'session/subscribed', sessionId: 's1', lastSeq: 1 } }
  197. const f2 = { type: 'server-request', rpcId: 'b', method: 'session/subscribed', payload: { type: 'session/subscribed', sessionId: 's2', lastSeq: 2 } }
  198. const wire = `: connected\n\ndata: ${JSON.stringify(f1)}\n\ndata: ${JSON.stringify(f2)}\n\n`
  199. const cuts = [5, 40, wire.indexOf('data: ', 40) + 3]
  200. const encoder = new TextEncoder()
  201. const doFetch = (): Promise<Response> => Promise.resolve(new Response(new ReadableStream<Uint8Array>({
  202. start(controller) {
  203. let prev = 0
  204. for (const cut of [...cuts, wire.length]) {
  205. controller.enqueue(encoder.encode(wire.slice(prev, cut)))
  206. prev = cut
  207. }
  208. controller.close()
  209. },
  210. }), { status: 200 }))
  211. const chopped = new InProcessApiClient({ fetch: doFetch })
  212. const seen: string[] = []
  213. for await (const envelope of chopped.events.mux({}, new AbortController().signal)) {
  214. seen.push((envelope.payload as { sessionId: string }).sessionId)
  215. expect(envelope.rpcId).toBe(seen.length === 1 ? 'a' : 'b')
  216. }
  217. expect(seen).toEqual(['s1', 's2'])
  218. })
  219. it('emits a stream/error frame then closes when the impl throws mid-stream', async () => {
  220. const api = scriptedApi({
  221. events: {
  222. async *host(request): AsyncGenerator<RpcRequest<HostFrame>> {
  223. yield { rpcId: RpcId(`p-${request.rpcId}`), payload: { type: 'host/session-added', sessionId: sid('s1') } }
  224. throw new Error('impl died mid-stream')
  225. },
  226. },
  227. })
  228. const seen: HostFrame[] = []
  229. for await (const envelope of client(api).events.host({}, new AbortController().signal)) {
  230. seen.push(envelope.payload)
  231. }
  232. expect(seen.map(f => f.type)).toEqual(['host/session-added', 'stream/error'])
  233. const last = seen.at(-1)
  234. if (last?.type === 'stream/error') expect(last.error.message).toMatch(/impl died mid-stream/)
  235. })
  236. it('drops a malformed SSE frame and keeps the stream alive (S→C two-level parse)', async () => {
  237. const good = { type: 'server-request', rpcId: 'g1', method: 'session/subscribed', payload: { type: 'session/subscribed', sessionId: 's1', lastSeq: 1 } }
  238. const badEnvelope = { type: 'server-response', rpcId: 'x' } // wrong quadrant for a stream
  239. const badFrame = { type: 'server-request', rpcId: 'b1', method: 'nope', payload: { type: 'no/such-frame' } }
  240. const wire = [
  241. 'data: {oops', // not JSON
  242. `data: ${JSON.stringify(badEnvelope)}`,
  243. `data: ${JSON.stringify(badFrame)}`,
  244. `data: ${JSON.stringify(good)}`,
  245. ].map(l => `${l}\n\n`).join('')
  246. const doFetch = (): Promise<Response> => Promise.resolve(new Response(new ReadableStream<Uint8Array>({
  247. start(controller) {
  248. controller.enqueue(new TextEncoder().encode(wire))
  249. controller.close()
  250. },
  251. }), { status: 200 }))
  252. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  253. try {
  254. const seen: MuxFrame[] = []
  255. for await (const envelope of new InProcessApiClient({ fetch: doFetch }).events.mux({}, new AbortController().signal)) {
  256. seen.push(envelope.payload)
  257. }
  258. // The three corrupt frames are reported and skipped; the good one still arrives.
  259. expect(seen).toEqual([{ type: 'session/subscribed', sessionId: 's1', lastSeq: 1 }])
  260. expect(errorSpy.mock.calls.length).toBe(3)
  261. } finally {
  262. errorSpy.mockRestore()
  263. }
  264. })
  265. it('fires onOpen once headers are in, before the first frame, and not on transport failure', async () => {
  266. const api = scriptedApi({
  267. events: {
  268. async *mux(request): AsyncGenerator<RpcRequest<MuxFrame>> {
  269. yield { rpcId: RpcId(`p-${request.rpcId}`), payload: { type: 'session/subscribed', sessionId: sid('s1'), lastSeq: 0 } }
  270. },
  271. },
  272. })
  273. const order: string[] = []
  274. const iterator = client(api).events.mux({}, new AbortController().signal, () => order.push('open'))
  275. expect(order).toEqual([]) // lazy generator: no fetch (and no onOpen) before iteration
  276. for await (const _ of iterator) order.push('frame')
  277. expect(order).toEqual(['open', 'frame'])
  278. // Transport failure path: onOpen must not fire.
  279. const failing = new InProcessApiClient({ fetch: () => Promise.resolve(new Response('down', { status: 503 })) })
  280. const failOrder: string[] = []
  281. await expect((async () => {
  282. for await (const _ of failing.events.mux({}, new AbortController().signal, () => failOrder.push('open'))) { /* unreachable */ }
  283. })()).rejects.toThrow(/transport failure/)
  284. expect(failOrder).toEqual([])
  285. })
  286. it('stops consuming when the caller aborts', async () => {
  287. let implSawAbort = false
  288. const api = scriptedApi({
  289. events: {
  290. async *mux(_request, signal): AsyncGenerator<RpcRequest<MuxFrame>> {
  291. try {
  292. let n = 0
  293. while (true) {
  294. yield { rpcId: RpcId(`p${n}`), payload: { type: 'session/subscribed', sessionId: sid('s1'), lastSeq: n++ } }
  295. await new Promise(resolve => setTimeout(resolve, 5))
  296. if (signal.aborted) return
  297. }
  298. } finally {
  299. implSawAbort = true
  300. }
  301. },
  302. },
  303. })
  304. const abort = new AbortController()
  305. let count = 0
  306. // In-process abort ends the stream (impl returns on signal.aborted); over a real
  307. // network fetch the same abort surfaces as a rejection — both stop the loop.
  308. await (async () => {
  309. for await (const _ of client(api).events.mux({}, abort.signal)) {
  310. if (++count === 2) abort.abort()
  311. }
  312. })().catch(() => undefined)
  313. expect(count).toBe(2)
  314. // Generator teardown may lag the abort by a microtask; poll briefly.
  315. await vi.waitFor(() => { expect(implSawAbort).toBe(true) })
  316. })
  317. })
  318. describe('respond path', () => {
  319. it('round-trips a client-response to a receipt', async () => {
  320. const seen: unknown[] = []
  321. const api = scriptedApi({
  322. respond: (message) => {
  323. seen.push(message)
  324. return Promise.resolve({ accepted: true as const })
  325. },
  326. })
  327. const receipt = await client(api).respond({ type: 'client-response', rpcId: RpcId('req-1'), result: { ok: true, value: { behavior: 'allow' } } })
  328. expect(receipt).toEqual({ accepted: true })
  329. expect(seen).toEqual([{ type: 'client-response', rpcId: 'req-1', result: { ok: true, value: { behavior: 'allow' } } }])
  330. })
  331. it('returns bad-response for a malformed client-response without reaching the impl', async () => {
  332. const respond = vi.fn()
  333. const handler = toFetchHandler(scriptedApi({ respond }))
  334. const response = await handler.fetch('http://dsh.internal/api/respond', { method: 'POST', body: JSON.stringify({ type: 'client-response' }) })
  335. expect(await response.json()).toEqual({ accepted: false, reason: 'bad-response' })
  336. expect(respond).not.toHaveBeenCalled()
  337. })
  338. })
  339. describe('envelope tap', () => {
  340. it('delivers one microtask batch of full forms per unary call', async () => {
  341. const api = scriptedApi()
  342. const tapped = client(api)
  343. const batches: (readonly RpcMessage[])[] = []
  344. tapped.subscribeEnvelopes(batch => batches.push(batch))
  345. await tapped.sessions.list({})
  346. await vi.waitFor(() => { expect(batches.length).toBeGreaterThan(0) })
  347. const all = batches.flat()
  348. expect(all.map(m => m.type)).toEqual(['client-request', 'server-response'])
  349. expect(all[0]?.rpcId).toBe(all[1]?.rpcId)
  350. })
  351. it('isolates a throwing listener and keeps serving the call', async () => {
  352. const api = scriptedApi()
  353. const tapped = client(api)
  354. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  355. try {
  356. const good: string[] = []
  357. tapped.subscribeEnvelopes(() => { throw new Error('listener bug') })
  358. tapped.subscribeEnvelopes(batch => good.push(...batch.map(m => m.type)))
  359. const response = await tapped.sessions.list({})
  360. expect(response.result.ok).toBe(true)
  361. await vi.waitFor(() => { expect(good).toContain('server-response') })
  362. } finally {
  363. errorSpy.mockRestore()
  364. }
  365. })
  366. it('buffers nothing with zero subscribers and unsubscribes cleanly', async () => {
  367. const api = scriptedApi()
  368. const tapped = client(api)
  369. await tapped.sessions.list({}) // no subscribers: must not accumulate
  370. const batches: (readonly RpcMessage[])[] = []
  371. const unsubscribe = tapped.subscribeEnvelopes(batch => batches.push(batch))
  372. unsubscribe()
  373. await tapped.sessions.list({})
  374. await new Promise(resolve => setTimeout(resolve, 0))
  375. expect(batches).toEqual([])
  376. })
  377. })