api-proxy-subagents.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363
  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: { source: { kind: string; rpcId: RpcId }; signal: AbortSignal },
  51. ) => options.followupError === undefined
  52. ? Promise.resolve('message-1')
  53. : Promise.reject(options.followupError))
  54. const interrupt = vi.fn((
  55. _targetSessionId: SessionId,
  56. _authority: { kind: 'user'; parentSessionId: SessionId },
  57. ) => {
  58. if (options.interruptError !== undefined) throw options.interruptError
  59. })
  60. const childHeader = {
  61. version: 0, id: CHILD, createdAt: 1, cwd: '/proj', parentSession: options.historyParent ?? PARENT,
  62. } satisfies SessionHeader
  63. const childEvents = [
  64. { type: 'user/message', seq: 0, time: 1, data: { content: [{ type: 'text', text: 'work' }], source: { kind: 'user' } } },
  65. ] as unknown as SessionEvent[]
  66. const inspect = vi.fn(() => Promise.resolve({ meta: childHeader, events: childEvents }))
  67. const liveBlock = { values: {}, asOfSeq: 3 }
  68. const coldBlock = { values: {}, asOfSeq: 0 }
  69. const snapshot = vi.fn(() => {
  70. if (options.projectionsThrow === true) throw new Error('hostile unit')
  71. return liveBlock
  72. })
  73. const restore = vi.fn(() => {
  74. if (options.projectionsThrow === true) throw new Error('hostile unit')
  75. return { snapshot: coldBlock }
  76. })
  77. const ctx = new Context()
  78. ctx.provide('agents', { get: getAgent })
  79. ctx.provide('subagents', { listChildren, followup, interrupt })
  80. ctx.provide('sessions', {
  81. get: (id: SessionId) => options.liveChild === true && id === CHILD
  82. ? { id: CHILD, header: childHeader, events: childEvents }
  83. : undefined,
  84. })
  85. ctx.provide('sessionPersistence', {
  86. list: () => Promise.resolve(options.storedChild === false ? [] : [childHeader]),
  87. inspect,
  88. locate: () => undefined,
  89. })
  90. // The gateway's own projection push feed subscribes at construction; the
  91. // no-op disposer keeps that feed quiet while these tests pin history reads.
  92. ctx.provide('sessionProjections', { snapshot, restore, onChanged: () => () => {} })
  93. ctx.provide('userInteraction', { registerProvider: () => () => {} })
  94. const api = createApiProxy(ctx, {
  95. defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp',
  96. })
  97. return { api, getAgent, listChildren, inspect, snapshot, restore, followup, interrupt, parent }
  98. }
  99. describe('subagent gateway', () => {
  100. it('lists the complete catalog and reports exact live-parent availability', async () => {
  101. const { api, listChildren } = bench({ parentLive: false, entries: [
  102. {
  103. kind: 'child', id: CHILD, mode: 'continuable', label: 'worker',
  104. activity: 'inactive', hasChildren: true,
  105. },
  106. {
  107. kind: 'child', id: sid('one-shot'), mode: 'one-shot',
  108. activity: 'inactive', hasChildren: false,
  109. },
  110. { kind: 'diagnostic', id: sid('bad'), reason: 'corrupt' },
  111. ] })
  112. const response = await api.subagents.list(request({ parentSessionId: PARENT }))
  113. expect(response.rpcId).toBe('subagent-rpc')
  114. expect(response.result).toMatchObject({
  115. ok: true,
  116. value: {
  117. parentAvailable: false,
  118. entries: [
  119. { kind: 'child', mode: 'continuable' },
  120. { kind: 'child', mode: 'one-shot' },
  121. { kind: 'diagnostic' },
  122. ],
  123. },
  124. })
  125. expect(listChildren).toHaveBeenCalledWith(PARENT, undefined)
  126. })
  127. it('derives catalog activity from the live child Agent rather than Session residency', async () => {
  128. const residentIdle = bench({ childStatus: 'idle', entries: [{
  129. kind: 'child', id: CHILD, mode: 'continuable', label: 'worker',
  130. activity: 'running', hasChildren: false,
  131. }] })
  132. expect((await residentIdle.api.subagents.list(request({ parentSessionId: PARENT }))).result)
  133. .toMatchObject({ ok: true, value: { entries: [{ activity: 'inactive' }] } })
  134. const running = bench({ childStatus: 'running' })
  135. expect((await running.api.subagents.list(request({ parentSessionId: PARENT }))).result)
  136. .toMatchObject({ ok: true, value: { entries: [{ activity: 'running' }] } })
  137. })
  138. it('reads a healthy direct child without looking up or activating any Agent', async () => {
  139. const { api, getAgent, inspect, restore } = bench()
  140. const response = await api.subagents.history(request({
  141. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', maxMessages: 10,
  142. }))
  143. expect(response.result).toMatchObject({
  144. ok: true,
  145. value: { hasMore: false, events: [{ event: { type: 'user/message', seq: 0 } }] },
  146. })
  147. expect(inspect).toHaveBeenCalledWith(CHILD)
  148. expect(restore).toHaveBeenCalledTimes(1)
  149. expect(getAgent).not.toHaveBeenCalled()
  150. })
  151. it('serves a live child from the in-memory snapshot and the watermark projections', async () => {
  152. const { api, inspect, snapshot, restore } = bench({ liveChild: true })
  153. const response = await api.subagents.history(request({
  154. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  155. }))
  156. expect(response.result).toMatchObject({
  157. ok: true,
  158. value: { hasMore: false, projections: { asOfSeq: 3 } },
  159. })
  160. expect(snapshot).toHaveBeenCalledTimes(1)
  161. expect(restore).not.toHaveBeenCalled()
  162. expect(inspect).not.toHaveBeenCalled()
  163. })
  164. it('serves the page without projections when a hostile unit breaks the fold', async () => {
  165. const cold = bench({ projectionsThrow: true })
  166. const coldResponse = await cold.api.subagents.history(request({
  167. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  168. }))
  169. expect(coldResponse.result).toMatchObject({
  170. ok: true,
  171. value: { hasMore: false, events: [{ event: { type: 'user/message', seq: 0 } }] },
  172. })
  173. if (coldResponse.result.ok) expect('projections' in coldResponse.result.value).toBe(false)
  174. const live = bench({ projectionsThrow: true, liveChild: true })
  175. const liveResponse = await live.api.subagents.history(request({
  176. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  177. }))
  178. expect(liveResponse.result).toMatchObject({
  179. ok: true,
  180. value: { hasMore: false, events: [{ event: { type: 'user/message', seq: 0 } }] },
  181. })
  182. if (liveResponse.result.ok) expect('projections' in liveResponse.result.value).toBe(false)
  183. expect(live.snapshot).toHaveBeenCalledTimes(1)
  184. })
  185. it('reads one-shot history and rejects an address with the wrong mode', async () => {
  186. const oneShot = {
  187. kind: 'child', id: CHILD, mode: 'one-shot', label: 'batch',
  188. activity: 'inactive', hasChildren: false,
  189. }
  190. const { api, inspect } = bench({ entries: [oneShot] })
  191. expect((await api.subagents.history(request({
  192. parentSessionId: PARENT, childSessionId: CHILD, mode: 'one-shot',
  193. }))).result).toMatchObject({ ok: true })
  194. expect((await api.subagents.history(request({
  195. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  196. }))).result).toMatchObject({ ok: false, error: { code: 'subagent-not-found' } })
  197. expect(inspect).toHaveBeenCalledTimes(1)
  198. })
  199. it('rejects a diagnostic address before reading history', async () => {
  200. const { api, inspect } = bench({ entries: [
  201. { kind: 'diagnostic', id: CHILD, reason: 'unsupported' },
  202. ] })
  203. const response = await api.subagents.history(request({
  204. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  205. }))
  206. expect(response.result).toMatchObject({
  207. ok: false,
  208. error: {
  209. code: 'subagent-catalog-diagnostic',
  210. details: { parentSessionId: PARENT, childSessionId: CHILD, reason: 'unsupported' },
  211. },
  212. })
  213. expect(inspect).not.toHaveBeenCalled()
  214. })
  215. it('maps the missing projections capability to one wire face on list, history, and prompt', async () => {
  216. const listError = () => new SubagentError(
  217. 'listing subagents requires the sessionProjections registry (load @deepseek-ai/dsh-session-projection)',
  218. 'SUBAGENT_CONTROL_PROJECTIONS_UNAVAILABLE',
  219. )
  220. const expected = {
  221. code: 'internal',
  222. message: 'subagent catalog is unavailable: this deployment does not mount the sessionProjections registry (load @deepseek-ai/dsh-session-projection)',
  223. }
  224. const list = bench({ listError: listError() })
  225. expect((await list.api.subagents.list(request({ parentSessionId: PARENT }))).result)
  226. .toMatchObject({ ok: false, error: expected })
  227. const history = bench({ listError: listError() })
  228. expect((await history.api.subagents.history(request({
  229. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  230. }))).result).toMatchObject({ ok: false, error: expected })
  231. expect(history.inspect).not.toHaveBeenCalled()
  232. const prompt = bench({ listError: listError() })
  233. expect((await prompt.api.subagents.prompt(request({
  234. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  235. }), new AbortController().signal)).result).toMatchObject({ ok: false, error: expected })
  236. expect(prompt.followup).not.toHaveBeenCalled()
  237. })
  238. it('routes human content through the exact live parent with rpc attribution', async () => {
  239. const { api, parent, followup } = bench()
  240. const content = [{ type: 'text' as const, text: '继续' }]
  241. const signal = new AbortController().signal
  242. const response = await api.subagents.prompt(request({
  243. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content,
  244. }), signal)
  245. expect(response.result).toMatchObject({
  246. ok: true, value: { messageId: 'message-1' },
  247. })
  248. expect(followup).toHaveBeenCalledWith(
  249. parent,
  250. CHILD,
  251. content,
  252. { source: { kind: 'user', rpcId: RpcId('subagent-rpc') }, signal },
  253. )
  254. })
  255. it('fails before delivery when the parent is absent and maps continuation failures', async () => {
  256. const absent = bench({ parentLive: false })
  257. expect((await absent.api.subagents.prompt(request({
  258. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  259. }), new AbortController().signal)).result).toMatchObject({
  260. ok: false, error: { code: 'subagent-parent-unavailable' },
  261. })
  262. expect(absent.listChildren).not.toHaveBeenCalled()
  263. const failed = bench({ followupError: new SubagentError('draining', 'DRAINING') })
  264. expect((await failed.api.subagents.prompt(request({
  265. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  266. }), new AbortController().signal)).result).toMatchObject({
  267. ok: false, error: { code: 'subagent-delivery-unavailable' },
  268. })
  269. })
  270. it('maps history disappearance and hides unexpected backend details', async () => {
  271. const disappeared = bench({ storedChild: false })
  272. expect((await disappeared.api.subagents.history(request({
  273. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable',
  274. }))).result).toMatchObject({
  275. ok: false,
  276. error: {
  277. code: 'subagent-not-found',
  278. message: 'subagent disappeared during history read',
  279. details: { parentSessionId: PARENT, childSessionId: CHILD },
  280. },
  281. })
  282. const catalog = bench({ listError: new Error('secret descriptor') })
  283. expect((await catalog.api.subagents.list(request({
  284. parentSessionId: PARENT,
  285. }))).result).toMatchObject({
  286. ok: false,
  287. error: { code: 'internal', message: 'subagent catalog read failed' },
  288. })
  289. const prompt = bench({ followupError: new Error('secret provider') })
  290. expect((await prompt.api.subagents.prompt(request({
  291. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable', content: [],
  292. }), new AbortController().signal)).result).toMatchObject({
  293. ok: false,
  294. error: { code: 'internal', message: 'subagent prompt failed' },
  295. })
  296. })
  297. it('interrupts through the core primitive alone while the parent Agent is offline', async () => {
  298. const { api, interrupt, getAgent, listChildren, inspect } = bench({ parentLive: false })
  299. const response = await api.subagents.interrupt(request({
  300. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable' as const,
  301. }))
  302. expect(response.rpcId).toBe('subagent-rpc')
  303. expect(response.result).toEqual({ ok: true, value: { accepted: true } })
  304. expect(interrupt).toHaveBeenCalledExactlyOnceWith(CHILD, { kind: 'user', parentSessionId: PARENT })
  305. // No parent-registry, catalog, or history dependency: this is what keeps a
  306. // live child interruptible after its parent Agent went offline.
  307. expect(getAgent).not.toHaveBeenCalled()
  308. expect(listChildren).not.toHaveBeenCalled()
  309. expect(inspect).not.toHaveBeenCalled()
  310. })
  311. it('maps interrupt authorization rejection without touching other services', async () => {
  312. const { api, listChildren } = bench({
  313. interruptError: new SubagentError('secret lineage', 'UNAUTHORIZED'),
  314. })
  315. const response = await api.subagents.interrupt(request({
  316. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable' as const,
  317. }))
  318. expect(response.result).toEqual({
  319. ok: false,
  320. error: {
  321. code: 'subagent-unauthorized',
  322. message: 'subagent does not belong to this parent',
  323. details: { childSessionId: CHILD },
  324. },
  325. })
  326. expect(listChildren).not.toHaveBeenCalled()
  327. })
  328. it('hides unexpected interrupt failures behind the internal code', async () => {
  329. const { api } = bench({ interruptError: new Error('secret activation state') })
  330. const response = await api.subagents.interrupt(request({
  331. parentSessionId: PARENT, childSessionId: CHILD, mode: 'continuable' as const,
  332. }))
  333. expect(response.result).toEqual({
  334. ok: false,
  335. error: { code: 'internal', message: 'subagent interrupt failed', details: {} },
  336. })
  337. })
  338. })