client-handler.spec.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442
  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, and rpcId discipline with no network or browser. Each case
  5. * 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, 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. subagents?: Partial<ApiProxy['subagents']>
  18. host?: Partial<ApiProxy['host']>
  19. skills?: Partial<ApiProxy['skills']>
  20. agentPresets?: Partial<ApiProxy['agentPresets']>
  21. settings?: Partial<ApiProxy['settings']>
  22. credentials?: Partial<ApiProxy['credentials']>
  23. llm?: Partial<ApiProxy['llm']>
  24. } = {}): ApiProxy {
  25. const err = <T>(r: RpcRequest<unknown>): Promise<RpcResponse<T>> =>
  26. Promise.resolve({ rpcId: r.rpcId, result: { ok: false, error: { code: 'internal' as const, message: 'stub', details: {} } } })
  27. return {
  28. subagents: {
  29. list: r => ok(r, { entries: [], parentAvailable: false }),
  30. prompt: r => ok(r, { messageId: 'message-1' as never }),
  31. interrupt: r => ok(r, { accepted: true as const }),
  32. ...overrides.subagents,
  33. },
  34. host: {
  35. describe: r => ok(r, {
  36. version: '0-test', cwd: '/t', attachedSessions: 0, home: '/h', canOpenPath: true,
  37. }),
  38. pickDirectory: r => ok(r, { path: null }),
  39. listDirectory: r => ok(r, { path: '/t', home: '/t', crumbs: [], entries: [], truncated: false }),
  40. createDirectory: r => ok(r, { path: '/t/new' }),
  41. openPath: r => ok(r, { opened: true as const }),
  42. ...overrides.host,
  43. },
  44. skills: { list: r => ok(r, { skills: [] }), ...overrides.skills },
  45. agentPresets: {
  46. openDocument: r => ok(r, { opened: true as const }),
  47. ...overrides.agentPresets,
  48. },
  49. settings: {
  50. describe: r => ok(r, { writable: true, hasDocument: false, namespaces: [] }),
  51. openDocument: r => ok(r, { opened: true as const }),
  52. update: err,
  53. replace: err,
  54. mutate: err,
  55. ...overrides.settings,
  56. },
  57. credentials: {
  58. describe: r => ok(r, { credentials: {} }),
  59. set: err,
  60. unset: err,
  61. ...overrides.credentials,
  62. },
  63. llm: {
  64. providers: r => ok(r, { providers: [] }),
  65. models: r => ok(r, {
  66. default: { provider: 'test', model: 'test' },
  67. routableProviders: [],
  68. groups: [],
  69. failures: [],
  70. }),
  71. discoverModels: err,
  72. ...overrides.llm,
  73. },
  74. downloads: { sessionLog: async () => new Response('stub', { status: 404 }) },
  75. }
  76. }
  77. function client(api: ApiProxy, timeoutMs?: number): InProcessApiClient {
  78. return new InProcessApiClient(toFetchHandler(api), timeoutMs)
  79. }
  80. /** Wrap one scripted method to record its invocation into `seen` before responding. */
  81. function recorderInto(seen: { method: string; payload: unknown }[]) {
  82. return <P, V>(method: string, respond: (r: RpcRequest<P>) => Promise<RpcResponse<V>>) =>
  83. (r: RpcRequest<P>): Promise<RpcResponse<V>> => {
  84. seen.push({ method, payload: r.payload })
  85. return respond(r)
  86. }
  87. }
  88. describe('unary round trip', () => {
  89. it('carries payload out and value back through the full wire form', async () => {
  90. let seen: RpcRequest<{}> | undefined
  91. const api = scriptedApi({
  92. host: {
  93. describe: (request) => {
  94. seen = request
  95. return ok(request, { version: '0-test', cwd: '/t', attachedSessions: 0, home: '/h', canOpenPath: true })
  96. },
  97. },
  98. })
  99. const response = await client(api).host.describe({})
  100. expect(seen?.payload).toEqual({})
  101. expect(seen?.rpcId).toBeTruthy()
  102. expect(response.rpcId).toBe(seen?.rpcId)
  103. expect(response.result).toMatchObject({ ok: true, value: { version: '0-test' } })
  104. })
  105. it('routes the agent-preset document opener through the wire', async () => {
  106. const opened = await client(scriptedApi()).agentPresets.openDocument({ agentPreset: 'mine' })
  107. expect(opened.result).toEqual({ ok: true, value: { opened: true } })
  108. })
  109. it('passes business errors through as 200 + err result, not a throw', async () => {
  110. const api = scriptedApi({
  111. host: {
  112. describe: request => Promise.resolve({
  113. rpcId: request.rpcId,
  114. result: { ok: false, error: { code: 'internal', message: 'nope', details: {} } },
  115. }),
  116. },
  117. })
  118. const response = await client(api).host.describe({})
  119. expect(response.result).toEqual({ ok: false, error: { code: 'internal', message: 'nope', details: {} } })
  120. })
  121. it('throws on rpcId echo mismatch', async () => {
  122. const api = scriptedApi({
  123. host: {
  124. describe: () => Promise.resolve({
  125. rpcId: RpcId('forged'),
  126. result: { ok: true, value: { version: '0-test', cwd: '/t', attachedSessions: 0, home: '/h', canOpenPath: true } },
  127. }),
  128. },
  129. })
  130. await expect(client(api).host.describe({})).rejects.toThrow(/rpcId mismatch/)
  131. })
  132. it('round-trips subagent.interrupt and rejects a one-shot or incomplete address', async () => {
  133. const interrupt = vi.fn((r: RpcRequest<unknown>) => ok(r, { accepted: true as const }))
  134. const api = scriptedApi({ subagents: { interrupt } })
  135. const c = client(api)
  136. const accepted = await c.subagents.interrupt({
  137. parentSessionId: sid('parent'), childSessionId: sid('child'), mode: 'continuable',
  138. })
  139. expect(accepted.result).toEqual({ ok: true, value: { accepted: true } })
  140. expect(interrupt).toHaveBeenCalledTimes(1)
  141. // The wire schema owns the mode fence: a one-shot address never reaches the impl.
  142. const oneShot = await c.subagents.interrupt({
  143. parentSessionId: sid('parent'), childSessionId: sid('child'), mode: 'one-shot',
  144. } as never)
  145. expect(oneShot.result.ok).toBe(false)
  146. if (!oneShot.result.ok) expect(oneShot.result.error.code).toBe('bad-request')
  147. const incomplete = await c.subagents.interrupt({
  148. parentSessionId: sid('parent'), mode: 'continuable',
  149. } as never)
  150. expect(incomplete.result.ok).toBe(false)
  151. if (!incomplete.result.ok) expect(incomplete.result.error.code).toBe('bad-request')
  152. expect(interrupt).toHaveBeenCalledTimes(1)
  153. })
  154. it('rejects a method/path mismatch as bad-request', async () => {
  155. const handler = toFetchHandler(scriptedApi())
  156. const body = { type: 'client-request', rpcId: 'r1', method: 'host.describe', payload: {} }
  157. const response = await handler.fetch('http://dsh.internal/api/skill.list', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify(body) })
  158. expect(response.status).toBe(200)
  159. const parsed = await response.json() as { result: { ok: boolean; error?: { code: string; message: string } } }
  160. expect(parsed.result.ok).toBe(false)
  161. expect(parsed.result.error?.code).toBe('bad-request')
  162. expect(parsed.result.error?.message).toMatch(/does not match path/)
  163. })
  164. it('rejects a malformed envelope as bad-request, salvaging the rpcId or falling back to the sentinel', async () => {
  165. const handler = toFetchHandler(scriptedApi())
  166. // No salvageable rpcId → the fixed invalid-request sentinel keeps the response a valid ServerResponse.
  167. const noId = await handler.fetch('http://dsh.internal/api/host.describe', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ nonsense: true }) })
  168. expect(noId.status).toBe(200)
  169. const noIdParsed = await noId.json() as { rpcId: string; result: { ok: boolean } }
  170. expect(noIdParsed.result.ok).toBe(false)
  171. expect(noIdParsed.rpcId).toBe('invalid-request')
  172. // A string rpcId in the otherwise-bad body is salvaged for correlation.
  173. const withId = await handler.fetch('http://dsh.internal/api/host.describe', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ rpcId: 'salvage-me', nonsense: true }) })
  174. const withIdParsed = await withId.json() as { rpcId: string; result: { ok: boolean } }
  175. expect(withIdParsed.result.ok).toBe(false)
  176. expect(withIdParsed.rpcId).toBe('salvage-me')
  177. })
  178. it('maps carrier failures to HTTP statuses and the client throws transport failure', async () => {
  179. const handler = toFetchHandler(scriptedApi())
  180. // Unknown method → 404.
  181. const notFound = await handler.fetch('http://dsh.internal/api/no.such', { method: 'POST', headers: { 'content-type': 'application/json' }, body: '{}' })
  182. expect(notFound.status).toBe(404)
  183. // Non-JSON body → 400.
  184. const badBody = await handler.fetch('http://dsh.internal/api/host.describe', { method: 'POST', headers: { 'content-type': 'application/json' }, body: '{oops' })
  185. expect(badBody.status).toBe(400)
  186. // Impl crash → 500, and through the client that is a throw, not an err result.
  187. const crashing = scriptedApi({ host: { describe: () => { throw new Error('impl exploded') } } })
  188. await expect(client(crashing).host.describe({})).rejects.toThrow(/transport failure .*500/)
  189. })
  190. it('rejects non-JSON media types before executing anything (cross-site simple-request fence)', async () => {
  191. const describe = vi.fn((request: RpcRequest<{}>) => ok(request, {
  192. version: '0-test', cwd: '/t', attachedSessions: 0, home: '/h', canOpenPath: true,
  193. }))
  194. const handler = toFetchHandler(scriptedApi({ host: { describe } }))
  195. const body = JSON.stringify({ type: 'client-request', rpcId: 'r1', method: 'host.describe', payload: {} })
  196. // A "simple" browser POST (text/plain — sent with no CORS preflight) is
  197. // refused at the carrier before the impl runs.
  198. const plain = await handler.fetch('http://dsh.internal/api/host.describe', { method: 'POST', headers: { 'content-type': 'text/plain' }, body })
  199. expect(plain.status).toBe(415)
  200. // A string body with no explicit header defaults to text/plain — same fence.
  201. const unlabelled = await handler.fetch('http://dsh.internal/api/host.describe', { method: 'POST', body })
  202. expect(unlabelled.status).toBe(415)
  203. expect(describe).not.toHaveBeenCalled()
  204. // Media-type parameters pass: the fence checks the type, not the exact string.
  205. const charset = await handler.fetch('http://dsh.internal/api/host.describe', { method: 'POST', headers: { 'content-type': 'application/json; charset=utf-8' }, body })
  206. expect(charset.status).toBe(200)
  207. expect(describe).toHaveBeenCalledTimes(1)
  208. })
  209. it('rejects when the transport never resolves within timeoutMs', async () => {
  210. // AbortSignal.timeout is immune to fake timers; a short real timeout keeps this fast.
  211. const never = new InProcessApiClient({
  212. fetch: (_i: RequestInfo | URL, init?: RequestInit) => new Promise<Response>((_resolve, reject) => {
  213. init?.signal?.addEventListener('abort', () => { reject(new Error('aborted by timeout')) })
  214. }),
  215. }, 25)
  216. await expect(never.host.describe({})).rejects.toThrow()
  217. })
  218. it('aborts a unary call through the caller-supplied external signal', async () => {
  219. // Real-fetch semantics: on abort the rejection is the signal's reason, and the abort
  220. // works even when the transport ignores the signal entirely (hung impl).
  221. const gate = new AbortController()
  222. const hung = new InProcessApiClient({ fetch: () => new Promise<Response>(() => {}) }, 60_000)
  223. const call = hung.host.describe({}, gate.signal)
  224. gate.abort(new Error('externally aborted'))
  225. await expect(call).rejects.toThrow(/externally aborted/)
  226. })
  227. it('rejects an already-aborted signal before touching the transport, mapping a string reason to an Error', async () => {
  228. let touched = false
  229. const c = new InProcessApiClient({
  230. fetch: () => {
  231. touched = true
  232. return Promise.resolve(new Response('{}'))
  233. },
  234. }, 60_000)
  235. const gate = new AbortController()
  236. gate.abort('gone before start')
  237. await expect(c.host.describe({}, gate.signal)).rejects.toThrow('gone before start')
  238. expect(touched).toBe(false)
  239. })
  240. it('maps a non-Error, non-string abort reason to the default AbortError message', async () => {
  241. const gate = new AbortController()
  242. const hung = new InProcessApiClient({ fetch: () => new Promise<Response>(() => {}) }, 60_000)
  243. const call = hung.host.describe({}, gate.signal)
  244. gate.abort(42)
  245. await expect(call).rejects.toThrow('This operation was aborted')
  246. })
  247. it('passes a signal-less doFetch straight through to the handler', async () => {
  248. class Probe extends InProcessApiClient {
  249. direct(url: URL): Promise<Response> {
  250. return this.doFetch(url)
  251. }
  252. }
  253. const probe = new Probe({ fetch: () => Promise.resolve(new Response('raw')) })
  254. const response = await probe.direct(new URL('http://dsh.internal/probe'))
  255. expect(await response.text()).toBe('raw')
  256. })
  257. it('throws on an S→C ok value that fails the method value schema (second-level parse)', async () => {
  258. // Impl echoes rpcId but returns a wrong-shaped value: envelope parse passes, value parse must reject.
  259. const api = scriptedApi({
  260. host: { describe: request => Promise.resolve({ rpcId: request.rpcId, result: { ok: true, value: { version: 1 } } }) as never },
  261. })
  262. await expect(client(api).host.describe({})).rejects.toThrow()
  263. })
  264. })
  265. describe('envelope tap', () => {
  266. it('delivers one microtask batch of full forms per unary call', async () => {
  267. const api = scriptedApi()
  268. const tapped = client(api)
  269. const batches: (readonly RpcMessage[])[] = []
  270. tapped.subscribeEnvelopes(batch => batches.push(batch))
  271. await tapped.host.describe({})
  272. await vi.waitFor(() => { expect(batches.length).toBeGreaterThan(0) })
  273. const all = batches.flat()
  274. expect(all.map(m => m.type)).toEqual(['client-request', 'server-response'])
  275. expect(all[0]?.rpcId).toBe(all[1]?.rpcId)
  276. })
  277. it('isolates a throwing listener and keeps serving the call', async () => {
  278. const api = scriptedApi()
  279. const tapped = client(api)
  280. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  281. try {
  282. const good: string[] = []
  283. tapped.subscribeEnvelopes(() => { throw new Error('listener bug') })
  284. tapped.subscribeEnvelopes(batch => good.push(...batch.map(m => m.type)))
  285. const response = await tapped.host.describe({})
  286. expect(response.result.ok).toBe(true)
  287. await vi.waitFor(() => { expect(good).toContain('server-response') })
  288. } finally {
  289. errorSpy.mockRestore()
  290. }
  291. })
  292. it('buffers nothing with zero subscribers and unsubscribes cleanly', async () => {
  293. const api = scriptedApi()
  294. const tapped = client(api)
  295. await tapped.host.describe({}) // no subscribers: must not accumulate
  296. const batches: (readonly RpcMessage[])[] = []
  297. const unsubscribe = tapped.subscribeEnvelopes(batch => batches.push(batch))
  298. unsubscribe()
  299. await tapped.host.describe({})
  300. await new Promise(resolve => setTimeout(resolve, 0))
  301. expect(batches).toEqual([])
  302. })
  303. })
  304. describe('config unary surface', () => {
  305. it('round-trips every settings/credentials/llm method with its own payload and value shape', async () => {
  306. const seen: { method: string; payload: unknown }[] = []
  307. const record = recorderInto(seen)
  308. const view = {
  309. ns: 'llm-deepseek',
  310. schema: { uid: 1, refs: { 1: { type: 'object' } } },
  311. value: { baseURL: 'https://next' },
  312. user: { baseURL: 'https://next' },
  313. applies: 'live' as const,
  314. secrets: [{ path: ['apiKey'], set: true }],
  315. revision: 0,
  316. }
  317. const providerRow = {
  318. provider: 'openai',
  319. displayName: 'openai',
  320. settingsNs: 'llm-pi-ai',
  321. settingsPath: ['providers', 'openai'],
  322. active: false,
  323. }
  324. const group = { id: 'deepseek-official', name: 'DeepSeek', models: [{ id: 'deepseek-v4-flash', name: 'Flash' }] }
  325. const api = scriptedApi({
  326. settings: {
  327. describe: record('settings.describe', r => ok(r, { writable: true, hasDocument: false, namespaces: [view] })),
  328. openDocument: record('settings.openDocument', r => ok(r, { opened: true as const })),
  329. update: record('settings.update', r => ok(r, view)),
  330. replace: record('settings.replace', r => ok(r, view)),
  331. mutate: record('settings.mutate', r => ok(r, view)),
  332. },
  333. credentials: {
  334. describe: record('credentials.describe', r => ok(r, { credentials: { OPENAI_API_KEY: { configured: true, source: 'file', writable: true } } })),
  335. set: record('credentials.set', r => ok(r, {})),
  336. unset: record('credentials.unset', r => ok(r, {})),
  337. },
  338. llm: {
  339. providers: record('llm.providers', r => ok(r, { providers: [providerRow] })),
  340. models: record('llm.models', r => ok(r, {
  341. default: { provider: 'deepseek-official', model: 'deepseek-v4-flash' },
  342. routableProviders: ['deepseek-official'],
  343. groups: [group],
  344. failures: [],
  345. })),
  346. discoverModels: record('llm.discoverModels', r => ok(r, { models: [{ id: 'acme-large', contextWindow: 65536 }] })),
  347. },
  348. })
  349. const c = client(api)
  350. const described = await c.settings.describe({})
  351. expect(described.result).toEqual({ ok: true, value: { writable: true, hasDocument: false, namespaces: [view] } })
  352. expect((await c.settings.openDocument({})).result).toEqual({ ok: true, value: { opened: true } })
  353. const updated = await c.settings.update({ ns: 'llm-deepseek', patch: { baseURL: 'https://next' } })
  354. expect(updated.result).toEqual({ ok: true, value: view })
  355. const replaced = await c.settings.replace({ ns: 'llm-deepseek', section: {} })
  356. expect(replaced.result).toEqual({ ok: true, value: view })
  357. const mutated = await c.settings.mutate({
  358. ns: 'llm-deepseek',
  359. ops: [{ op: 'unset', path: ['baseURL'] }],
  360. expectedRevision: 0,
  361. })
  362. expect(mutated.result).toEqual({ ok: true, value: view })
  363. const creds = await c.credentials.describe({ refs: ['OPENAI_API_KEY'] })
  364. expect(creds.result).toEqual({ ok: true, value: { credentials: { OPENAI_API_KEY: { configured: true, source: 'file', writable: true } } } })
  365. expect((await c.credentials.set({ ref: 'OPENAI_API_KEY', value: 'sk-x' })).result).toEqual({ ok: true, value: {} })
  366. expect((await c.credentials.unset({ ref: 'OPENAI_API_KEY' })).result).toEqual({ ok: true, value: {} })
  367. const providers = await c.llm.providers({})
  368. expect(providers.result).toEqual({ ok: true, value: { providers: [providerRow] } })
  369. const models = await c.llm.models({})
  370. expect(models.result).toEqual({
  371. ok: true,
  372. value: {
  373. default: { provider: 'deepseek-official', model: 'deepseek-v4-flash' },
  374. routableProviders: ['deepseek-official'],
  375. groups: [group],
  376. failures: [],
  377. },
  378. })
  379. const discovered = await c.llm.discoverModels({
  380. settingsNs: 'llm-pi-ai',
  381. baseURL: 'https://gateway.acme.example/v1',
  382. api: 'openai-completions',
  383. apiKey: 'probe-key',
  384. })
  385. expect(discovered.result).toEqual({ ok: true, value: { models: [{ id: 'acme-large', contextWindow: 65536 }] } })
  386. expect(seen.map(call => call.method)).toEqual([
  387. 'settings.describe', 'settings.openDocument', 'settings.update', 'settings.replace', 'settings.mutate',
  388. 'credentials.describe', 'credentials.set', 'credentials.unset',
  389. 'llm.providers', 'llm.models', 'llm.discoverModels',
  390. ])
  391. expect(seen[2]?.payload).toEqual({ ns: 'llm-deepseek', patch: { baseURL: 'https://next' } })
  392. expect(seen[4]?.payload)
  393. .toEqual({ ns: 'llm-deepseek', ops: [{ op: 'unset', path: ['baseURL'] }], expectedRevision: 0 })
  394. expect(seen[6]?.payload).toEqual({ ref: 'OPENAI_API_KEY', value: 'sk-x' })
  395. // The draft crosses whole, credential included: the host needs it for this
  396. // one interrogation and stores none of it.
  397. expect(seen[10]?.payload).toEqual({
  398. settingsNs: 'llm-pi-ai',
  399. baseURL: 'https://gateway.acme.example/v1',
  400. api: 'openai-completions',
  401. apiKey: 'probe-key',
  402. })
  403. })
  404. it('rejects an invalid credential reference name at the carrier boundary', async () => {
  405. const api = scriptedApi()
  406. const response = await client(api).credentials.set({ ref: 'not a var', value: 'x' })
  407. expect(response.result.ok).toBe(false)
  408. if (response.result.ok) throw new Error('unreachable')
  409. expect(response.result.error.code).toBe('bad-request')
  410. })
  411. })