control-approval.host.spec.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450
  1. import { Context } from '@deepseek-ai/cordis'
  2. import AgentRegistry from '@deepseek-ai/dsh-agent'
  3. import type { Agent } from '@deepseek-ai/dsh-agent'
  4. import SessionStore from '@deepseek-ai/dsh-session'
  5. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  6. import ApprovalService from '@deepseek-ai/dsh-user-approval'
  7. import type { ApprovalRequestId } from '@deepseek-ai/dsh-user-approval'
  8. import UserQuestionService from '@deepseek-ai/dsh-user-questions'
  9. import { describe, expect, it, vi } from 'vitest'
  10. import { SessionControlController } from '../src/control.ts'
  11. import type {
  12. SessionApprovalRequest,
  13. SessionControlFrame,
  14. SessionInteractionId,
  15. SessionRespondRequest,
  16. } from '../src/types.ts'
  17. interface ControlCapture {
  18. readonly frames: SessionControlFrame[]
  19. waitFor(type: SessionControlFrame['type']): Promise<SessionControlFrame>
  20. }
  21. async function harness(): Promise<{ ctx: Context; control: SessionControlController }> {
  22. const ctx = new Context()
  23. await ctx.plugin(SessionStore)
  24. await ctx.plugin(SystemPrompt, { persona: '' })
  25. await ctx.plugin(UserQuestionService)
  26. await ctx.plugin(AgentRegistry)
  27. await ctx.plugin(ApprovalService)
  28. const control = new SessionControlController(ctx)
  29. await new Promise(resolve => setTimeout(resolve, 0))
  30. return { ctx, control }
  31. }
  32. function agentOf(ctx: Context): Agent {
  33. const session = ctx.sessions.create()
  34. session.append('turn/start', { turn: 1 })
  35. return { session } as unknown as Agent
  36. }
  37. function openControl(control: SessionControlController, abort: AbortController): ControlCapture {
  38. const frames: SessionControlFrame[] = []
  39. const waiters: {
  40. type: SessionControlFrame['type']
  41. resolve(frame: SessionControlFrame): void
  42. }[] = []
  43. void (async () => {
  44. for await (const frame of control.control(abort.signal)) {
  45. frames.push(frame)
  46. for (let index = waiters.length - 1; index >= 0; index--) {
  47. const waiter = waiters[index] as (typeof waiters)[number]
  48. if (waiter.type !== frame.type) continue
  49. waiters.splice(index, 1)
  50. waiter.resolve(frame)
  51. }
  52. }
  53. })()
  54. return {
  55. frames,
  56. waitFor: (type) => {
  57. const found = frames.find(frame => frame.type === type)
  58. if (found !== undefined) return Promise.resolve(found)
  59. return new Promise((resolve) => { waiters.push({ type, resolve }) })
  60. },
  61. }
  62. }
  63. function requestedOf(frame: SessionControlFrame): SessionApprovalRequest {
  64. if (frame.type !== 'approval/requested') {
  65. throw new Error(`expected approval/requested, got ${frame.type}`)
  66. }
  67. return frame
  68. }
  69. async function waitForCount(
  70. stream: ControlCapture,
  71. type: SessionControlFrame['type'],
  72. count: number,
  73. ): Promise<void> {
  74. for (let index = 0; index < 200 && stream.frames.filter(frame => frame.type === type).length < count; index++) {
  75. await new Promise(resolve => setTimeout(resolve, 5))
  76. }
  77. expect(stream.frames.filter(frame => frame.type === type).length).toBeGreaterThanOrEqual(count)
  78. }
  79. function answer(
  80. interactionId: SessionInteractionId,
  81. sessionId: unknown,
  82. approvalId: ApprovalRequestId,
  83. outcome: 'allowed-once' | 'rejected',
  84. ): SessionRespondRequest {
  85. return {
  86. interactionId,
  87. result: { ok: true, value: { sessionId, approvalId, outcome } },
  88. } as SessionRespondRequest
  89. }
  90. describe('approval pending registry', () => {
  91. it('round-trips ask through requested, response, outcome, and resolved frames', async () => {
  92. const { ctx, control } = await harness()
  93. const abort = new AbortController()
  94. const stream = openControl(control, abort)
  95. const agent = agentOf(ctx)
  96. const asked = ctx.approval.request({ agent, toolName: 'bash', reason: 'sandbox escalation' })
  97. const requested = requestedOf(await stream.waitFor('approval/requested'))
  98. expect(requested).toMatchObject({
  99. toolName: 'bash',
  100. reason: 'sandbox escalation',
  101. sessionId: agent.session.id,
  102. })
  103. expect(control.respond(answer(
  104. requested.interactionId,
  105. requested.sessionId,
  106. requested.approvalId,
  107. 'allowed-once',
  108. ))).toEqual({ accepted: true })
  109. await expect(asked).resolves.toBe('allowed-once')
  110. const resolved = await stream.waitFor('approval/resolved')
  111. expect(resolved).toMatchObject({ approvalId: requested.approvalId, outcome: 'allowed-once' })
  112. expect(control.respond(answer(
  113. requested.interactionId,
  114. requested.sessionId,
  115. requested.approvalId,
  116. 'rejected',
  117. ))).toEqual({ accepted: false, reason: 'not-pending' })
  118. abort.abort()
  119. })
  120. it('replays one pending request with the same interaction id in a new baseline', async () => {
  121. const { ctx, control } = await harness()
  122. const firstAbort = new AbortController()
  123. const first = openControl(control, firstAbort)
  124. const agent = agentOf(ctx)
  125. const asked = ctx.approval.request({ agent, toolName: 'write' })
  126. const requested = requestedOf(await first.waitFor('approval/requested'))
  127. firstAbort.abort()
  128. const secondAbort = new AbortController()
  129. const second = openControl(control, secondAbort)
  130. const baseline = await second.waitFor('baseline')
  131. if (baseline.type !== 'baseline') throw new Error('expected baseline')
  132. const replayed = baseline.value.approvals[0]
  133. expect(replayed?.interactionId).toBe(requested.interactionId)
  134. expect(replayed?.approvalId).toBe(requested.approvalId)
  135. expect(control.respond(answer(
  136. requested.interactionId,
  137. requested.sessionId,
  138. requested.approvalId,
  139. 'rejected',
  140. ))).toEqual({ accepted: true })
  141. await expect(asked).resolves.toBe('rejected')
  142. secondAbort.abort()
  143. })
  144. it('rejects malformed and mismatched answers, and reports unknown interactions', async () => {
  145. const { ctx, control } = await harness()
  146. const abort = new AbortController()
  147. const stream = openControl(control, abort)
  148. const agent = agentOf(ctx)
  149. void ctx.approval.request({ agent, toolName: 'bash' })
  150. const requested = requestedOf(await stream.waitFor('approval/requested'))
  151. expect(control.respond(answer(
  152. 'ghost' as SessionInteractionId,
  153. requested.sessionId,
  154. requested.approvalId,
  155. 'rejected',
  156. ))).toEqual({ accepted: false, reason: 'not-pending' })
  157. expect(control.respond({
  158. interactionId: requested.interactionId,
  159. result: { ok: false, error: { code: 'internal', message: 'x', details: {} } },
  160. })).toEqual({ accepted: false, reason: 'bad-response' })
  161. expect(control.respond(answer(
  162. requested.interactionId,
  163. requested.sessionId,
  164. 'other-approval' as ApprovalRequestId,
  165. 'rejected',
  166. ))).toEqual({ accepted: false, reason: 'bad-response' })
  167. expect(control.respond({
  168. interactionId: requested.interactionId,
  169. result: { ok: true, value: { nonsense: 1 } as never },
  170. })).toEqual({ accepted: false, reason: 'bad-response' })
  171. abort.abort()
  172. })
  173. it('withdraws an approval when its ask signal aborts', async () => {
  174. const { ctx, control } = await harness()
  175. const abort = new AbortController()
  176. const stream = openControl(control, abort)
  177. const agent = agentOf(ctx)
  178. const cancel = new AbortController()
  179. const asked = ctx.approval.request({ agent, toolName: 'bash', signal: cancel.signal })
  180. const requested = requestedOf(await stream.waitFor('approval/requested'))
  181. cancel.abort()
  182. await expect(asked).resolves.toBe('cancelled')
  183. expect(await stream.waitFor('approval/resolved')).toMatchObject({
  184. approvalId: requested.approvalId,
  185. outcome: 'cancelled',
  186. })
  187. expect(control.respond(answer(
  188. requested.interactionId,
  189. requested.sessionId,
  190. requested.approvalId,
  191. 'allowed-once',
  192. ))).toEqual({ accepted: false, reason: 'not-pending' })
  193. abort.abort()
  194. })
  195. it('settles a pre-aborted dispatch without publishing it', async () => {
  196. const { ctx, control } = await harness()
  197. const abort = new AbortController()
  198. const stream = openControl(control, abort)
  199. const session = ctx.sessions.create()
  200. session.append('turn/start', { turn: 1 })
  201. session.append('approval/asked', {
  202. id: 'pre-aborted' as ApprovalRequestId,
  203. toolName: 'bash',
  204. })
  205. const agent = { session } as unknown as Agent
  206. const cancelled = new AbortController()
  207. cancelled.abort()
  208. const outcome = await ctx.waterfall(
  209. 'approval/request',
  210. { agent, toolName: 'bash', signal: cancelled.signal },
  211. () => Promise.resolve('unavailable' as const),
  212. )
  213. expect(outcome).toBe('cancelled')
  214. const secondAbort = new AbortController()
  215. const second = openControl(control, secondAbort)
  216. await second.waitFor('baseline')
  217. expect(second.frames.some(frame => frame.type === 'approval/requested')).toBe(false)
  218. secondAbort.abort()
  219. abort.abort()
  220. void stream
  221. })
  222. it('settles pending approvals when the controller is disposed', async () => {
  223. const ctx = new Context()
  224. await ctx.plugin(SessionStore)
  225. await ctx.plugin(SystemPrompt, { persona: '' })
  226. await ctx.plugin(UserQuestionService)
  227. await ctx.plugin(AgentRegistry)
  228. await ctx.plugin(ApprovalService)
  229. let control!: SessionControlController
  230. const fiber = ctx.plugin(Object.assign((fiberCtx: Context) => {
  231. control = new SessionControlController(fiberCtx)
  232. }, { inject: ['sessions', 'agents', 'userQuestions', 'approval'] }))
  233. await fiber.await()
  234. const abort = new AbortController()
  235. const stream = openControl(control, abort)
  236. const asked = ctx.approval.request({ agent: agentOf(ctx), toolName: 'bash' })
  237. const requested = requestedOf(await stream.waitFor('approval/requested'))
  238. await fiber.dispose()
  239. await expect(asked).resolves.toBe('cancelled')
  240. expect(await stream.waitFor('approval/resolved')).toMatchObject({
  241. approvalId: requested.approvalId,
  242. outcome: 'cancelled',
  243. })
  244. abort.abort()
  245. })
  246. it('carries callId and ignores an abort after the answer settled', async () => {
  247. const { ctx, control } = await harness()
  248. const abort = new AbortController()
  249. const stream = openControl(control, abort)
  250. const agent = agentOf(ctx)
  251. const cancel = new AbortController()
  252. vi.spyOn(cancel.signal, 'removeEventListener').mockImplementation(() => {})
  253. const asked = ctx.approval.request({
  254. agent,
  255. toolName: 'bash',
  256. callId: 'call-9' as never,
  257. signal: cancel.signal,
  258. })
  259. const requested = requestedOf(await stream.waitFor('approval/requested'))
  260. expect(requested.callId).toBe('call-9')
  261. expect(control.respond(answer(
  262. requested.interactionId,
  263. requested.sessionId,
  264. requested.approvalId,
  265. 'allowed-once',
  266. ))).toEqual({ accepted: true })
  267. await expect(asked).resolves.toBe('allowed-once')
  268. cancel.abort()
  269. expect(stream.frames.filter(frame => frame.type === 'approval/resolved')).toHaveLength(1)
  270. abort.abort()
  271. })
  272. it('contains an abort that wins immediately after pending registration', async () => {
  273. const { ctx, control } = await harness()
  274. void control
  275. let reads = 0
  276. const signal = {
  277. get aborted() { return ++reads >= 3 },
  278. addEventListener: () => {},
  279. removeEventListener: () => {},
  280. } as unknown as AbortSignal
  281. await expect(ctx.approval.request({
  282. agent: agentOf(ctx),
  283. toolName: 'bash',
  284. signal,
  285. })).resolves.toBe('cancelled')
  286. })
  287. it('cancels only approvals owned by a disposed Session', async () => {
  288. const { ctx, control } = await harness()
  289. const abort = new AbortController()
  290. const stream = openControl(control, abort)
  291. const first = agentOf(ctx)
  292. const second = agentOf(ctx)
  293. const firstAsk = ctx.approval.request({ agent: first, toolName: 'first' })
  294. const secondAsk = ctx.approval.request({ agent: second, toolName: 'second' })
  295. await waitForCount(stream, 'approval/requested', 2)
  296. const requests = stream.frames.filter(
  297. (frame): frame is Extract<SessionControlFrame, { type: 'approval/requested' }> => (
  298. frame.type === 'approval/requested'
  299. ),
  300. )
  301. ctx.emit('session/disposed', first.session)
  302. await expect(firstAsk).resolves.toBe('cancelled')
  303. const remaining = requests.find(request => request.sessionId === second.session.id)
  304. if (remaining === undefined) throw new Error('missing second approval')
  305. expect(control.respond(answer(
  306. remaining.interactionId,
  307. remaining.sessionId,
  308. remaining.approvalId,
  309. 'allowed-once',
  310. ))).toEqual({ accepted: true })
  311. await expect(secondAsk).resolves.toBe('allowed-once')
  312. abort.abort()
  313. })
  314. it('pairs parallel asks by callId', async () => {
  315. const { ctx, control } = await harness()
  316. const abort = new AbortController()
  317. const stream = openControl(control, abort)
  318. const agent = agentOf(ctx)
  319. const askA = ctx.approval.request({ agent, toolName: 'bash', callId: 'call-a' as never })
  320. const askB = ctx.approval.request({ agent, toolName: 'bash', callId: 'call-b' as never })
  321. await waitForCount(stream, 'approval/requested', 2)
  322. const requests = stream.frames
  323. .filter((frame): frame is Extract<SessionControlFrame, { type: 'approval/requested' }> => (
  324. frame.type === 'approval/requested'
  325. ))
  326. const requestA = requests.find(frame => frame.callId === 'call-a')
  327. const requestB = requests.find(frame => frame.callId === 'call-b')
  328. if (requestA === undefined || requestB === undefined) throw new Error('missing parallel request')
  329. const askedIdByCall = new Map(agent.session.events
  330. .filter(event => event.type === 'approval/asked')
  331. .map(event => [String(event.data.callId), event.data.id]))
  332. expect(requestA.approvalId).toBe(askedIdByCall.get('call-a'))
  333. expect(requestB.approvalId).toBe(askedIdByCall.get('call-b'))
  334. expect(control.respond(answer(
  335. requestB.interactionId,
  336. requestB.sessionId,
  337. requestB.approvalId,
  338. 'rejected',
  339. ))).toEqual({ accepted: true })
  340. expect(control.respond(answer(
  341. requestA.interactionId,
  342. requestA.sessionId,
  343. requestA.approvalId,
  344. 'allowed-once',
  345. ))).toEqual({ accepted: true })
  346. await expect(askA).resolves.toBe('allowed-once')
  347. await expect(askB).resolves.toBe('rejected')
  348. abort.abort()
  349. })
  350. it('gives parallel callId-less asks distinct audit ids', async () => {
  351. const { ctx, control } = await harness()
  352. const abort = new AbortController()
  353. const stream = openControl(control, abort)
  354. const agent = agentOf(ctx)
  355. const askA = ctx.approval.request({ agent, toolName: 'alpha' })
  356. const askB = ctx.approval.request({ agent, toolName: 'beta' })
  357. await waitForCount(stream, 'approval/requested', 2)
  358. const requests = stream.frames
  359. .filter((frame): frame is Extract<SessionControlFrame, { type: 'approval/requested' }> => (
  360. frame.type === 'approval/requested'
  361. ))
  362. const requestA = requests.find(frame => frame.toolName === 'alpha')
  363. const requestB = requests.find(frame => frame.toolName === 'beta')
  364. if (requestA === undefined || requestB === undefined) throw new Error('missing parallel request')
  365. expect(requestA.approvalId).not.toBe(requestB.approvalId)
  366. expect(control.respond(answer(
  367. requestA.interactionId,
  368. requestA.sessionId,
  369. requestA.approvalId,
  370. 'allowed-once',
  371. ))).toEqual({ accepted: true })
  372. expect(control.respond(answer(
  373. requestB.interactionId,
  374. requestB.sessionId,
  375. requestB.approvalId,
  376. 'rejected',
  377. ))).toEqual({ accepted: true })
  378. await expect(askA).resolves.toBe('allowed-once')
  379. await expect(askB).resolves.toBe('rejected')
  380. abort.abort()
  381. })
  382. it('delegates a dispatch whose only asked candidate is already decided', async () => {
  383. const { ctx, control } = await harness()
  384. void control
  385. const session = ctx.sessions.create()
  386. session.append('turn/start', { turn: 1 })
  387. session.append('approval/asked', {
  388. id: 'stale-ask' as ApprovalRequestId,
  389. toolName: 'bash',
  390. })
  391. session.append('approval/decided', {
  392. id: 'stale-ask' as ApprovalRequestId,
  393. outcome: 'rejected',
  394. })
  395. const agent = { session } as unknown as Agent
  396. const outcome = await ctx.waterfall(
  397. 'approval/request',
  398. { agent, toolName: 'bash' },
  399. () => Promise.resolve('unavailable' as const),
  400. )
  401. expect(outcome).toBe('unavailable')
  402. })
  403. it('delegates an ask with no matching audit event', async () => {
  404. const { ctx, control } = await harness()
  405. void control
  406. const session = ctx.sessions.create()
  407. session.append('turn/start', { turn: 1 })
  408. const agent = { session } as unknown as Agent
  409. const outcome = await ctx.waterfall(
  410. 'approval/request',
  411. { agent, toolName: 'x' },
  412. () => Promise.resolve('unavailable' as const),
  413. )
  414. expect(outcome).toBe('unavailable')
  415. })
  416. })