api-proxy-subagents.spec.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import type { SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  4. import { SubagentError } from '@deepseek-ai/dsh-subagent'
  5. import { RpcId } from '../src/api/rpc.ts'
  6. import type { RpcRequest } from '../src/api/rpc.ts'
  7. import { createApiProxy } from '../src/api-proxy.ts'
  8. const sid = (value: string): SessionId => value as SessionId
  9. const PARENT = sid('parent')
  10. const CHILD = sid('child')
  11. function request<P>(payload: P): RpcRequest<P> {
  12. return { rpcId: RpcId('subagent-rpc'), payload }
  13. }
  14. function bench(options: {
  15. parentLive?: boolean
  16. childStatus?: 'idle' | 'running'
  17. entries?: object[]
  18. followupError?: Error
  19. interruptError?: Error
  20. listError?: Error
  21. /** Persistence forgets the child entirely (the vanished-mid-read race). */
  22. storedChild?: false
  23. /** Attach the child to the live session store instead of persistence only. */
  24. liveChild?: true
  25. /** Every registered projection unit throws on this child's payloads. */
  26. projectionsThrow?: true
  27. historyParent?: SessionId
  28. } = {}) {
  29. const parent = { id: PARENT }
  30. const child = options.childStatus === undefined
  31. ? undefined
  32. : { id: CHILD, status: options.childStatus }
  33. const getAgent = vi.fn((id: SessionId) => {
  34. if (options.parentLive !== false && id === PARENT) return parent
  35. if (id === CHILD) return child
  36. return undefined
  37. })
  38. const listChildren = vi.fn(() => options.listError === undefined
  39. ? Promise.resolve(options.entries ?? [
  40. {
  41. kind: 'child', id: CHILD, mode: 'continuable', label: 'worker',
  42. activity: 'inactive', hasChildren: false,
  43. },
  44. ])
  45. : Promise.reject(options.listError))
  46. const followup = vi.fn((
  47. _parent: unknown,
  48. _childId: SessionId,
  49. _content: unknown,
  50. _delivery: {
  51. source: { kind: string; rpcId: RpcId; clientTimeZone?: string }
  52. signal: AbortSignal
  53. },
  54. ) => options.followupError === undefined
  55. ? Promise.resolve('message-1')
  56. : Promise.reject(options.followupError))
  57. const interrupt = vi.fn((
  58. _targetSessionId: SessionId,
  59. _authority: { kind: 'user'; parentSessionId: SessionId },
  60. ) => {
  61. if (options.interruptError !== undefined) throw options.interruptError
  62. })
  63. const childHeader = {
  64. version: 0, id: CHILD, createdAt: 1, cwd: '/proj', parentSession: options.historyParent ?? PARENT,
  65. } satisfies SessionHeader
  66. const childEvents = [
  67. { type: 'user/message', seq: 0, time: 1, data: { content: [{ type: 'text', text: 'work' }], source: { kind: 'user' } } },
  68. ] as unknown as SessionEvent[]
  69. const inspect = vi.fn(() => Promise.resolve({ meta: childHeader, events: childEvents }))
  70. const liveBlock = { values: {}, asOfSeq: 3 }
  71. const coldBlock = { values: {}, asOfSeq: 0 }
  72. const snapshot = vi.fn(() => {
  73. if (options.projectionsThrow === true) throw new Error('hostile unit')
  74. return liveBlock
  75. })
  76. const restore = vi.fn(() => {
  77. if (options.projectionsThrow === true) throw new Error('hostile unit')
  78. return { snapshot: coldBlock }
  79. })
  80. const ctx = new Context()
  81. ctx.provide('agents', { get: getAgent })
  82. ctx.provide('subagents', { listChildren, followup, interrupt })
  83. ctx.provide('sessions', {
  84. get: (id: SessionId) => options.liveChild === true && id === CHILD
  85. ? { id: CHILD, header: childHeader, events: childEvents }
  86. : undefined,
  87. })
  88. ctx.provide('sessionPersistence', {
  89. list: () => Promise.resolve(options.storedChild === false ? [] : [childHeader]),
  90. inspect,
  91. locate: () => undefined,
  92. })
  93. // The gateway's own projection push feed subscribes at construction; the
  94. // no-op disposer keeps that feed quiet while these tests pin history reads.
  95. ctx.provide('sessionProjections', {
  96. snapshot,
  97. restore,
  98. onChanged: () => () => {},
  99. register: () => () => {},
  100. })
  101. ctx.provide('userQuestions', { registerProvider: () => () => {} })
  102. const api = createApiProxy(ctx, {
  103. defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp',
  104. })
  105. return { api, getAgent, listChildren, inspect, snapshot, restore, followup, interrupt, parent }
  106. }
  107. describe('subagent gateway', () => {
  108. it('lists the complete catalog and reports exact live-parent availability', async () => {
  109. const { api, listChildren } = bench({ parentLive: false, entries: [
  110. {
  111. kind: 'child', id: CHILD, mode: 'continuable', label: 'worker',
  112. activity: 'inactive', hasChildren: true,
  113. },
  114. {
  115. kind: 'child', id: sid('one-shot'), mode: 'one-shot',
  116. activity: 'inactive', hasChildren: false,
  117. },
  118. { kind: 'diagnostic', id: sid('bad'), reason: 'corrupt' },
  119. ] })
  120. const response = await api.subagents.list(request({ parentSessionId: PARENT }))
  121. expect(response.rpcId).toBe('subagent-rpc')
  122. expect(response.result).toMatchObject({
  123. ok: true,
  124. value: {
  125. parentAvailable: false,
  126. entries: [
  127. { kind: 'child', mode: 'continuable' },
  128. { kind: 'child', mode: 'one-shot' },
  129. { kind: 'diagnostic' },
  130. ],
  131. },
  132. })
  133. expect(listChildren).toHaveBeenCalledWith(PARENT, undefined)
  134. })
  135. it('derives catalog activity from the live child Agent rather than Session residency', async () => {
  136. const residentIdle = bench({ childStatus: 'idle', entries: [{
  137. kind: 'child', id: CHILD, mode: 'continuable', label: 'worker',
  138. activity: 'running', hasChildren: false,
  139. }] })
  140. expect((await residentIdle.api.subagents.list(request({ parentSessionId: PARENT }))).result)
  141. .toMatchObject({ ok: true, value: { entries: [{ activity: 'inactive' }] } })
  142. const running = bench({ childStatus: 'running' })
  143. expect((await running.api.subagents.list(request({ parentSessionId: PARENT }))).result)
  144. .toMatchObject({ ok: true, value: { entries: [{ activity: 'running' }] } })
  145. })
  146. it('reads a healthy direct child without looking up or activating any Agent', async () => {
  147. const { api, getAgent, inspect, restore } = bench()
  148. const response = await api.subagents.history(request({
  149. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', maxMessages: 10,
  150. }))
  151. expect(response.result).toMatchObject({
  152. ok: true,
  153. value: { hasMore: false, events: [{ event: { type: 'user/message', seq: 0 } }] },
  154. })
  155. expect(inspect).toHaveBeenCalledWith(CHILD)
  156. expect(restore).toHaveBeenCalledTimes(1)
  157. expect(getAgent).not.toHaveBeenCalled()
  158. })
  159. it('serves a live child from the in-memory snapshot and the watermark projections', async () => {
  160. const { api, inspect, snapshot, restore } = bench({ liveChild: true })
  161. const response = await api.subagents.history(request({
  162. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  163. }))
  164. expect(response.result).toMatchObject({
  165. ok: true,
  166. value: { hasMore: false, projections: { asOfSeq: 3 } },
  167. })
  168. expect(snapshot).toHaveBeenCalledTimes(1)
  169. expect(restore).not.toHaveBeenCalled()
  170. expect(inspect).not.toHaveBeenCalled()
  171. })
  172. it('serves the page without projections when a hostile unit breaks the fold', async () => {
  173. const cold = bench({ projectionsThrow: true })
  174. const coldResponse = await cold.api.subagents.history(request({
  175. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  176. }))
  177. expect(coldResponse.result).toMatchObject({
  178. ok: true,
  179. value: { hasMore: false, events: [{ event: { type: 'user/message', seq: 0 } }] },
  180. })
  181. if (coldResponse.result.ok) expect('projections' in coldResponse.result.value).toBe(false)
  182. const live = bench({ projectionsThrow: true, liveChild: true })
  183. const liveResponse = await live.api.subagents.history(request({
  184. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  185. }))
  186. expect(liveResponse.result).toMatchObject({
  187. ok: true,
  188. value: { hasMore: false, events: [{ event: { type: 'user/message', seq: 0 } }] },
  189. })
  190. if (liveResponse.result.ok) expect('projections' in liveResponse.result.value).toBe(false)
  191. expect(live.snapshot).toHaveBeenCalledTimes(1)
  192. })
  193. it('reads one-shot history and rejects an address with the wrong mode', async () => {
  194. const oneShot = {
  195. kind: 'child', id: CHILD, mode: 'one-shot', label: 'batch',
  196. activity: 'inactive', hasChildren: false,
  197. }
  198. const { api, inspect } = bench({ entries: [oneShot] })
  199. expect((await api.subagents.history(request({
  200. parentSessionId: PARENT, childSessionId: CHILD, mode: 'one-shot',
  201. }))).result).toMatchObject({ ok: true })
  202. expect((await api.subagents.history(request({
  203. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  204. }))).result).toMatchObject({ ok: false, error: { code: 'subagent-not-found' } })
  205. expect(inspect).toHaveBeenCalledTimes(1)
  206. })
  207. it('rejects a diagnostic address before reading history', async () => {
  208. const { api, inspect } = bench({ entries: [
  209. { kind: 'diagnostic', id: CHILD, reason: 'unsupported' },
  210. ] })
  211. const response = await api.subagents.history(request({
  212. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  213. }))
  214. expect(response.result).toMatchObject({
  215. ok: false,
  216. error: {
  217. code: 'subagent-catalog-diagnostic',
  218. details: { parentSessionId: PARENT, childSessionId: CHILD, reason: 'unsupported' },
  219. },
  220. })
  221. expect(inspect).not.toHaveBeenCalled()
  222. })
  223. it('maps the missing projections capability to one wire face on list, history, and prompt', async () => {
  224. const listError = () => new SubagentError(
  225. 'listing subagents requires the sessionProjections registry (load @deepseek-ai/dsh-session-projection)',
  226. 'SUBAGENT_CONTROL_PROJECTIONS_UNAVAILABLE',
  227. )
  228. const expected = {
  229. code: 'internal',
  230. message: 'subagent catalog is unavailable: this deployment does not mount the sessionProjections registry (load @deepseek-ai/dsh-session-projection)',
  231. }
  232. const list = bench({ listError: listError() })
  233. expect((await list.api.subagents.list(request({ parentSessionId: PARENT }))).result)
  234. .toMatchObject({ ok: false, error: expected })
  235. const history = bench({ listError: listError() })
  236. expect((await history.api.subagents.history(request({
  237. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  238. }))).result).toMatchObject({ ok: false, error: expected })
  239. expect(history.inspect).not.toHaveBeenCalled()
  240. const prompt = bench({ listError: listError() })
  241. expect((await prompt.api.subagents.prompt(request({
  242. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  243. }), new AbortController().signal)).result).toMatchObject({ ok: false, error: expected })
  244. expect(prompt.followup).not.toHaveBeenCalled()
  245. })
  246. it('routes human content through the exact live parent with rpc attribution', async () => {
  247. const { api, parent, followup } = bench()
  248. const content = [{ type: 'text' as const, text: '继续' }]
  249. const signal = new AbortController().signal
  250. const response = await api.subagents.prompt(request({
  251. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content,
  252. }), signal)
  253. expect(response.result).toMatchObject({
  254. ok: true, value: { messageId: 'message-1' },
  255. })
  256. expect(followup).toHaveBeenCalledWith(
  257. parent,
  258. CHILD,
  259. content,
  260. { source: { kind: 'user', rpcId: RpcId('subagent-rpc') }, signal },
  261. )
  262. })
  263. it('canonicalizes browser-zone provenance before delivering a child prompt', async () => {
  264. const { api, parent, followup } = bench()
  265. const alias = 'US/Pacific'
  266. const canonical = new Intl.DateTimeFormat('en-US', { timeZone: alias })
  267. .resolvedOptions().timeZone
  268. const content = [{ type: 'text' as const, text: 'continue locally' }]
  269. const signal = new AbortController().signal
  270. await expect(api.subagents.prompt(request({
  271. parentSessionId: PARENT,
  272. childSessionId: CHILD,
  273. mode: 'continuable',
  274. content,
  275. clientTimeZone: alias,
  276. }), signal)).resolves.toMatchObject({ result: { ok: true } })
  277. expect(followup).toHaveBeenCalledWith(parent, CHILD, content, {
  278. source: { kind: 'user', rpcId: RpcId('subagent-rpc'), clientTimeZone: canonical },
  279. signal,
  280. })
  281. const invalid = await api.subagents.prompt(request({
  282. parentSessionId: PARENT,
  283. childSessionId: CHILD,
  284. mode: 'continuable',
  285. content,
  286. clientTimeZone: 'Not/A_Real_Zone',
  287. }), signal)
  288. expect(invalid.result).toEqual({
  289. ok: false,
  290. error: {
  291. code: 'invalid-time-zone',
  292. message: 'clientTimeZone must be UTC or a valid IANA Area/Location name',
  293. details: { value: 'Not/A_Real_Zone' },
  294. },
  295. })
  296. expect(followup).toHaveBeenCalledOnce()
  297. })
  298. it('fails before delivery when the parent is absent and maps continuation failures', async () => {
  299. const absent = bench({ parentLive: false })
  300. expect((await absent.api.subagents.prompt(request({
  301. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  302. }), new AbortController().signal)).result).toMatchObject({
  303. ok: false, error: { code: 'subagent-parent-unavailable' },
  304. })
  305. expect(absent.listChildren).not.toHaveBeenCalled()
  306. const failed = bench({ followupError: new SubagentError('draining', 'DRAINING') })
  307. expect((await failed.api.subagents.prompt(request({
  308. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  309. }), new AbortController().signal)).result).toMatchObject({
  310. ok: false, error: { code: 'subagent-delivery-unavailable' },
  311. })
  312. })
  313. it('maps history disappearance and hides unexpected backend details', async () => {
  314. const disappeared = bench({ storedChild: false })
  315. expect((await disappeared.api.subagents.history(request({
  316. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  317. }))).result).toMatchObject({
  318. ok: false,
  319. error: {
  320. code: 'subagent-not-found',
  321. message: 'subagent disappeared during history read',
  322. details: { parentSessionId: PARENT, childSessionId: CHILD },
  323. },
  324. })
  325. const catalog = bench({ listError: new Error('secret descriptor') })
  326. expect((await catalog.api.subagents.list(request({
  327. parentSessionId: PARENT,
  328. }))).result).toMatchObject({
  329. ok: false,
  330. error: { code: 'internal', message: 'subagent catalog read failed' },
  331. })
  332. const prompt = bench({ followupError: new Error('secret provider') })
  333. expect((await prompt.api.subagents.prompt(request({
  334. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  335. }), new AbortController().signal)).result).toMatchObject({
  336. ok: false,
  337. error: { code: 'internal', message: 'subagent prompt failed' },
  338. })
  339. })
  340. it('interrupts through the core primitive alone while the parent Agent is offline', async () => {
  341. const { api, interrupt, getAgent, listChildren, inspect } = bench({ parentLive: false })
  342. const response = await api.subagents.interrupt(request({
  343. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable' as const,
  344. }))
  345. expect(response.rpcId).toBe('subagent-rpc')
  346. expect(response.result).toEqual({ ok: true, value: { accepted: true } })
  347. expect(interrupt).toHaveBeenCalledExactlyOnceWith(CHILD, { kind: 'user', parentSessionId: PARENT })
  348. // No parent-registry, catalog, or history dependency: this is what keeps a
  349. // live child interruptible after its parent Agent went offline.
  350. expect(getAgent).not.toHaveBeenCalled()
  351. expect(listChildren).not.toHaveBeenCalled()
  352. expect(inspect).not.toHaveBeenCalled()
  353. })
  354. it('maps interrupt authorization rejection without touching other services', async () => {
  355. const { api, listChildren } = bench({
  356. interruptError: new SubagentError('secret lineage', 'UNAUTHORIZED'),
  357. })
  358. const response = await api.subagents.interrupt(request({
  359. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable' as const,
  360. }))
  361. expect(response.result).toEqual({
  362. ok: false,
  363. error: {
  364. code: 'subagent-unauthorized',
  365. message: 'subagent does not belong to this parent',
  366. details: { childSessionId: CHILD },
  367. },
  368. })
  369. expect(listChildren).not.toHaveBeenCalled()
  370. })
  371. it('hides unexpected interrupt failures behind the internal code', async () => {
  372. const { api } = bench({ interruptError: new Error('secret activation state') })
  373. const response = await api.subagents.interrupt(request({
  374. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable' as const,
  375. }))
  376. expect(response.result).toEqual({
  377. ok: false,
  378. error: { code: 'internal', message: 'subagent interrupt failed', details: {} },
  379. })
  380. })
  381. })