interception.spec.ts 39 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import LlmRuntime, { createUserMessage, ToolCallId } from '@deepseek-ai/dsh-llm'
  4. import SessionStore, {
  5. SessionId,
  6. type SessionEvent,
  7. type TurnEndReason,
  8. type UserMessage,
  9. } from '@deepseek-ai/dsh-session'
  10. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  11. import ToolRuntime, { defineContentToolFixture, type PostToolDecision, type PreToolDecision } from '@deepseek-ai/dsh-tools'
  12. import AgentRegistry, {
  13. type Agent,
  14. type PreStepDecision,
  15. type SessionStartSource,
  16. } from '@deepseek-ai/dsh-agent'
  17. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  18. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  19. import { MockAdapter, textResponse, toolCallResponse } from './mock-adapter.ts'
  20. /**
  21. * The interception points introduced by the hooks taxonomy: `agent/pre-step`,
  22. * `agent/session-start`, `agent/turn-stopping`, and the
  23. * `tools/pre-execute` / `tools/post-execute`
  24. * split with `additionalContexts` buffering. These verify the canonical event
  25. * API a hook bridge (or a native plugin) programs against, WITHOUT any
  26. * external protocol — a native plugin uses the typed decisions directly.
  27. */
  28. async function harness(adapter: MockAdapter) {
  29. const ctx = new Context()
  30. await ctx.plugin(LlmRuntime)
  31. await ctx.plugin(SessionStore)
  32. await ctx.plugin(SessionProjectionRegistry)
  33. await ctx.plugin(SystemPrompt)
  34. await ctx.plugin(ToolRuntime)
  35. await ctx.plugin(AgentRegistry)
  36. await ctx.plugin(AgentLoop, { agents: [] })
  37. ctx.llm.registerAdapter(['mock'], adapter)
  38. return ctx
  39. }
  40. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  41. return new Promise((resolve) => {
  42. const dispose = ctx.on('agent/status', ({ agent: subject, status }) => {
  43. if (subject === agent && status === 'idle') {
  44. dispose()
  45. resolve()
  46. }
  47. })
  48. })
  49. }
  50. function send(agent: Agent, text: string) {
  51. agent.followup(createUserMessage({ content: [{ type: 'text', text }], source: { kind: 'user' } }))
  52. }
  53. function events(agent: Agent): readonly SessionEvent[] {
  54. return agent.session.snapshotEvents()
  55. }
  56. describe('agent/pre-step', () => {
  57. it('supports following-only, surface-only, and empty Agent send options', async () => {
  58. const adapter = new MockAdapter([textResponse('one'), textResponse('two'), textResponse('three')])
  59. const ctx = await harness(adapter)
  60. const agent = await ctx.agentLoop.create(SessionId('send-options'), { provider: 'mock', model: 'mock' })
  61. const first = createUserMessage({ content: [{ type: 'text', text: 'first' }], source: { kind: 'user' } })
  62. const context = createUserMessage({
  63. content: [{ type: 'text', text: 'following' }], source: { kind: 'plugin', plugin: 'test' },
  64. })
  65. let idle = waitForIdle(ctx, agent)
  66. agent.send(first, 'next-turn', true, { followingMessages: [context] })
  67. expect(() => agent.runMaintenance(() => Promise.resolve())).toThrow('already has active work')
  68. await idle
  69. expect(adapter.requests[0]!.messages).toMatchObject([
  70. { content: [{ type: 'text', text: 'first' }] },
  71. { content: [{ type: 'text', text: 'following' }] },
  72. ])
  73. const second = createUserMessage({ content: [{ type: 'text', text: 'second' }], source: { kind: 'user' } })
  74. idle = waitForIdle(ctx, agent)
  75. agent.send(second, 'next-turn', true, { surfaceIntent: { surfaceOp: 'append' } })
  76. await idle
  77. expect(events(agent).find(event => event.type === 'user/message' && event.data.id === second.id))
  78. .toMatchObject({ surfaceOp: 'append' })
  79. const third = createUserMessage({ content: [{ type: 'text', text: 'third' }], source: { kind: 'user' } })
  80. idle = waitForIdle(ctx, agent)
  81. agent.send(third, 'next-turn', true, { followingMessages: [] })
  82. await idle
  83. expect(adapter.requests).toHaveLength(3)
  84. })
  85. it('records a claimed message with its durable replacement intent after listener rewrites', async () => {
  86. const adapter = new MockAdapter([textResponse('old answer'), textResponse('new answer')])
  87. const ctx = await harness(adapter)
  88. const agent = await ctx.agentLoop.create(SessionId('surface-intent'), { provider: 'mock', model: 'mock' })
  89. send(agent, 'original')
  90. await waitForIdle(ctx, agent)
  91. const before = events(agent)
  92. const original = before.find(event => event.type === 'user/message')
  93. const oldAnswer = before.find(event => event.type === 'assistant/message')
  94. const turnStart = before.find(event => event.type === 'turn/start')
  95. const turnEnd = before.findLast(event => event.type === 'turn/end')
  96. if (original?.type !== 'user/message' || oldAnswer?.type !== 'assistant/message'
  97. || turnStart?.type !== 'turn/start' || turnEnd?.type !== 'turn/end') {
  98. throw new Error('fixture did not complete its first turn')
  99. }
  100. const edited = createUserMessage({
  101. content: [{ type: 'text', text: 'edited' }],
  102. source: { kind: 'user' },
  103. })
  104. const context = createUserMessage({
  105. content: [{ type: 'text', text: 'preserved context' }],
  106. source: { kind: 'plugin', plugin: 'test' },
  107. })
  108. const surfaceIntent = {
  109. surfaceOp: { op: 'replace' as const, start: original.seq, end: oldAnswer.seq },
  110. sourceEventSeqs: [original.seq, oldAnswer.seq],
  111. conversationOp: { op: 'replace' as const, start: turnStart.seq, end: turnEnd.seq },
  112. }
  113. ctx.on('agent/pre-step', async ({ surfaceIntents }, next) => {
  114. expect(surfaceIntents?.get(edited.id)).toEqual(surfaceIntent)
  115. const decision = await next()
  116. if (decision.kind === 'reject') return decision
  117. return {
  118. ...decision,
  119. messages: decision.messages.map(message => message.id === edited.id
  120. ? { ...message, content: [{ type: 'text', text: 'rewritten edit' }] }
  121. : message),
  122. }
  123. })
  124. const idle = waitForIdle(ctx, agent)
  125. agent.send(edited, 'next-turn', true, {
  126. position: 'front',
  127. followingMessages: [context],
  128. surfaceIntent,
  129. })
  130. await idle
  131. const replacement = events(agent).find(event =>
  132. event.type === 'user/message' && event.data.id === edited.id)
  133. expect(replacement).toMatchObject({
  134. type: 'user/message',
  135. data: { content: [{ type: 'text', text: 'rewritten edit' }] },
  136. surfaceOp: { op: 'replace', start: original.seq, end: oldAnswer.seq },
  137. sourceEventSeqs: [original.seq, oldAnswer.seq],
  138. conversationOp: { op: 'replace', start: turnStart.seq, end: turnEnd.seq },
  139. })
  140. expect(adapter.requests[1]!.messages).toMatchObject([
  141. { content: [{ type: 'text', text: 'rewritten edit' }] },
  142. { content: [{ type: 'text', text: 'preserved context' }] },
  143. ])
  144. expect(JSON.stringify(adapter.requests[1]!.messages)).not.toContain('original')
  145. expect(JSON.stringify(adapter.requests[1]!.messages)).not.toContain('old answer')
  146. })
  147. it('enter (default via next) records the user/message unchanged', async () => {
  148. const adapter = new MockAdapter([textResponse('ok')])
  149. const ctx = await harness(adapter)
  150. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  151. const seen: string[] = []
  152. ctx.on('agent/pre-step', async ({ messages }, next) => {
  153. seen.push(messages[0]!.content.map(b => (b.type === 'text' ? b.text : '')).join(''))
  154. return next()
  155. })
  156. send(agent, 'hello')
  157. await waitForIdle(ctx, agent)
  158. expect(seen).toEqual(['hello'])
  159. const userMsg = events(agent).find(e => e.type === 'user/message')
  160. expect(userMsg?.type === 'user/message' && userMsg.data.content).toEqual([{ type: 'text', text: 'hello' }])
  161. })
  162. it('reports the request coordinates for initial and tool-continuation prompts', async () => {
  163. const adapter = new MockAdapter([
  164. toolCallResponse('c1', 'echo', { text: 'hi' }),
  165. textResponse('done'),
  166. ])
  167. const ctx = await harness(adapter)
  168. ctx.tools.register(defineContentToolFixture({
  169. name: 'echo',
  170. description: 'echo',
  171. parameters: { text: { type: 'string', required: true } },
  172. execute: async ({ text }) => [{ type: 'text', text }],
  173. }))
  174. const agent = await ctx.agentLoop.create(SessionId('prompt-coordinates'), { provider: 'mock', model: 'mock' })
  175. const seen: Array<{ turn: number; step: number; messages: number }> = []
  176. ctx.on('agent/pre-step', async ({ messages, turn, step }, next) => {
  177. seen.push({ turn, step, messages: messages.length })
  178. return next()
  179. })
  180. send(agent, 'hello')
  181. await waitForIdle(ctx, agent)
  182. expect(seen).toEqual([
  183. { turn: 1, step: 1, messages: 1 },
  184. { turn: 1, step: 2, messages: 0 },
  185. ])
  186. })
  187. it('publishes frozen input without replacing its identity', async () => {
  188. const adapter = new MockAdapter([textResponse('ok')])
  189. const ctx = await harness(adapter)
  190. const agent = await ctx.agentLoop.create(SessionId('owned-input'), { provider: 'mock', model: 'mock' })
  191. const entered = Promise.withResolvers<undefined>()
  192. const decision = Promise.withResolvers<PreStepDecision>()
  193. const observed: UserMessage[] = []
  194. ctx.on('agent/pre-step', async ({ agent: subject, messages }) => {
  195. if (subject !== agent) return { kind: 'enter', messages }
  196. const message = messages[0]!
  197. expect(Object.isFrozen(message)).toBe(true)
  198. expect(Object.isFrozen(message.content)).toBe(true)
  199. expect(Object.isFrozen(message.content[0])).toBe(true)
  200. expect(Object.isFrozen(message.source)).toBe(true)
  201. expect(() => {
  202. const block = message.content[0]
  203. if (block?.type === 'text') block.text = 'listener mutation'
  204. }).toThrow()
  205. observed.push(message)
  206. entered.resolve(undefined)
  207. return decision.promise
  208. })
  209. const input: UserMessage = createUserMessage({
  210. content: [{ type: 'text', text: 'accepted text' }],
  211. source: { kind: 'plugin', plugin: 'accepted source' },
  212. })
  213. const idle = waitForIdle(ctx, agent)
  214. agent.followup(input)
  215. await entered.promise
  216. const block = input.content[0]
  217. expect(() => {
  218. if (block?.type === 'text') block.text = 'caller mutation'
  219. }).toThrow(TypeError)
  220. expect(() => {
  221. if (input.source.kind === 'plugin') input.source.plugin = 'caller mutation'
  222. }).toThrow(TypeError)
  223. decision.resolve({ kind: 'enter', messages: [input] })
  224. await idle
  225. expect(observed).toHaveLength(1)
  226. expect(observed[0]).not.toBe(input)
  227. expect(observed[0]).toMatchObject({
  228. content: [{ type: 'text', text: 'accepted text' }],
  229. source: { kind: 'plugin', plugin: 'accepted source' },
  230. })
  231. const userMsg = events(agent).find(event => event.type === 'user/message')
  232. expect(userMsg?.type === 'user/message' && userMsg.data).toEqual(input)
  233. })
  234. it('enter with content rewrites the prompt before it is recorded', async () => {
  235. const adapter = new MockAdapter([textResponse('ok')])
  236. const ctx = await harness(adapter)
  237. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  238. ctx.on('agent/pre-step', async ({ messages }): Promise<PreStepDecision> =>
  239. ({
  240. kind: 'enter',
  241. messages: [{ ...messages[0]!, content: [{ type: 'text', text: 'REWRITTEN' }] }],
  242. }))
  243. send(agent, 'original')
  244. await waitForIdle(ctx, agent)
  245. const userMsg = events(agent).find(e => e.type === 'user/message')
  246. expect(userMsg?.type === 'user/message' && userMsg.data.content).toEqual([{ type: 'text', text: 'REWRITTEN' }])
  247. // the rewritten prompt is what reached the model
  248. expect(JSON.stringify(adapter.requests[0]!.messages)).toContain('REWRITTEN')
  249. expect(JSON.stringify(adapter.requests[0]!.messages)).not.toContain('original')
  250. })
  251. it('enter with additional messages records separately sourced context in the turn', async () => {
  252. const adapter = new MockAdapter([textResponse('ok')])
  253. const ctx = await harness(adapter)
  254. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  255. ctx.on('agent/pre-step', async ({ messages }): Promise<PreStepDecision> =>
  256. ({
  257. kind: 'enter',
  258. messages: [...messages, createUserMessage({
  259. content: [{ type: 'text', text: '<system-reminder>extra ctx</system-reminder>' }],
  260. source: { kind: 'plugin', plugin: 'test' },
  261. })],
  262. }))
  263. send(agent, 'go')
  264. await waitForIdle(ctx, agent)
  265. const log = events(agent)
  266. const userMsg = log.find(e => e.type === 'user/message' && e.data.source.kind === 'user')
  267. const ctxMsg = log.find(e => e.type === 'user/message' && e.data.source.kind === 'plugin')
  268. expect(userMsg).toBeDefined()
  269. expect(ctxMsg?.type === 'user/message' && ctxMsg.data.content).toEqual([{ type: 'text', text: '<system-reminder>extra ctx</system-reminder>' }])
  270. expect(ctxMsg?.type === 'user/message' && ctxMsg.data.source).toEqual({ kind: 'plugin', plugin: 'test' })
  271. const sent = JSON.stringify(adapter.requests[0]!.messages)
  272. expect(sent).toContain('extra ctx')
  273. })
  274. it('does not open another step when a completed turn rewrites pending input to empty', async () => {
  275. const adapter = new MockAdapter([textResponse('done')])
  276. const ctx = await harness(adapter)
  277. const agent = await ctx.agentLoop.create(SessionId('empty-completed-continuation'), {
  278. provider: 'mock',
  279. model: 'mock',
  280. })
  281. ctx.on('agent/turn-stopping', ({ agent: subject }) => {
  282. subject.inject(createUserMessage({
  283. content: [{ type: 'text', text: 'pending context' }],
  284. source: { kind: 'plugin', plugin: 'test' },
  285. }))
  286. })
  287. ctx.on('agent/pre-step', async ({ step }, next) => {
  288. const decision = await next()
  289. return step === 1 || decision.kind === 'reject'
  290. ? decision
  291. : { kind: 'enter', messages: [] }
  292. })
  293. send(agent, 'finish once')
  294. await agent.whenIdle()
  295. expect(adapter.requests).toHaveLength(1)
  296. expect(events(agent).filter(event => event.type === 'step/start')).toHaveLength(1)
  297. })
  298. it('reject closes the claimed prompt turn without a step or model call', async () => {
  299. const adapter = new MockAdapter([textResponse('should not run')])
  300. const ctx = await harness(adapter)
  301. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  302. ctx.on('agent/pre-step', async (): Promise<PreStepDecision> => ({ kind: 'reject' }))
  303. const reasons: TurnEndReason[] = []
  304. ctx.on('session/event', (_s, event: SessionEvent) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  305. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'do something' }], source: { kind: 'user' } }))
  306. await agent.whenIdle()
  307. // the model was never called
  308. expect(adapter.requests).toHaveLength(0)
  309. const log = events(agent)
  310. expect(log.filter(e => e.type === 'turn/start' || e.type === 'turn/end').map(e => e.type))
  311. .toEqual(['turn/start', 'turn/end'])
  312. expect(log.some(e => e.type === 'user/message')).toBe(false)
  313. expect(log.some(e => e.type === 'step/start')).toBe(false)
  314. expect(reasons).toEqual([{ kind: 'blocked' }])
  315. })
  316. it('stages inject and steer during pre-step for the entered turn', async () => {
  317. const adapter = new MockAdapter([textResponse('ok')])
  318. const ctx = await harness(adapter)
  319. const agent = await ctx.agentLoop.create(SessionId('pre-step-outbox'), { provider: 'mock', model: 'mock' })
  320. const entered = Promise.withResolvers<undefined>()
  321. const decision = Promise.withResolvers<PreStepDecision>()
  322. let claimed: UserMessage[] = []
  323. let firstProposal = true
  324. ctx.on('agent/pre-step', async ({ messages }) => {
  325. if (!firstProposal) return { kind: 'enter', messages }
  326. firstProposal = false
  327. claimed = messages
  328. entered.resolve(undefined)
  329. return decision.promise
  330. })
  331. const idle = waitForIdle(ctx, agent)
  332. send(agent, 'entered prompt')
  333. await entered.promise
  334. expect(agent.status).toBe('running')
  335. expect(events(agent).some(event => event.type === 'turn/start')).toBe(true)
  336. agent.inject(createUserMessage({
  337. content: [{ type: 'text', text: 'attached context' }],
  338. source: { kind: 'plugin', plugin: 'test' },
  339. }))
  340. agent.steer(createUserMessage({ content: [{ type: 'text', text: 'pre-step steering' }], source: { kind: 'user' } }))
  341. expect(events(agent).some(event => event.type === 'user/message')).toBe(false)
  342. expect(agent.inbox.nextStep.map(message => message.content[0]))
  343. .toEqual([
  344. { type: 'text', text: 'attached context' },
  345. { type: 'text', text: 'pre-step steering' },
  346. ])
  347. decision.resolve({ kind: 'enter', messages: claimed })
  348. await idle
  349. expect(agent.inbox.hasPending).toBe(false)
  350. const staged = events(agent).filter(event =>
  351. event.type === 'turn/start' || event.type === 'user/message')
  352. expect(staged.map(event => event.type)).toEqual([
  353. 'turn/start',
  354. 'user/message',
  355. 'user/message',
  356. 'user/message',
  357. ])
  358. expect(staged[1]?.type === 'user/message' && staged[1].data.content)
  359. .toEqual([{ type: 'text', text: 'entered prompt' }])
  360. expect(staged[2]?.type === 'user/message' && staged[2].data.content)
  361. .toEqual([{ type: 'text', text: 'attached context' }])
  362. expect(staged[3]?.type === 'user/message' && staged[3].data.content)
  363. .toEqual([{ type: 'text', text: 'pre-step steering' }])
  364. const firstRequest = JSON.stringify(adapter.requests[0]?.messages)
  365. expect(firstRequest).toContain('entered prompt')
  366. expect(firstRequest).not.toContain('attached context')
  367. expect(firstRequest).not.toContain('pre-step steering')
  368. const nextRequest = JSON.stringify(adapter.requests[1]?.messages)
  369. expect(nextRequest).toContain('attached context')
  370. expect(nextRequest).toContain('pre-step steering')
  371. })
  372. it('preserves input staged after the blocked batch was claimed', async () => {
  373. const adapter = new MockAdapter([textResponse('retried')])
  374. const ctx = await harness(adapter)
  375. const agent = await ctx.agentLoop.create(SessionId('blocked-pre-step-outbox'), { provider: 'mock', model: 'mock' })
  376. const entered = Promise.withResolvers<undefined>()
  377. const decision = Promise.withResolvers<PreStepDecision>()
  378. const disposeBlock = ctx.on('agent/pre-step', async () => {
  379. entered.resolve(undefined)
  380. return decision.promise
  381. })
  382. const blockedIdle = waitForIdle(ctx, agent)
  383. send(agent, 'blocked prompt')
  384. await entered.promise
  385. agent.inject(createUserMessage({
  386. content: [{ type: 'text', text: 'staged context' }],
  387. source: { kind: 'plugin', plugin: 'test' },
  388. }))
  389. agent.steer(createUserMessage({ content: [{ type: 'text', text: 'staged steering' }], source: { kind: 'user' } }))
  390. decision.resolve({ kind: 'reject' })
  391. await blockedIdle
  392. expect(agent.inbox.nextStep.map(message => message.content[0]))
  393. .toEqual([
  394. { type: 'text', text: 'staged context' },
  395. { type: 'text', text: 'staged steering' },
  396. ])
  397. expect(events(agent).filter(event => event.type === 'turn/start' || event.type === 'turn/end')
  398. .map(event => event.type)).toEqual(['turn/start', 'turn/end'])
  399. expect(adapter.requests).toEqual([])
  400. disposeBlock()
  401. send(agent, 'resume')
  402. await waitForIdle(ctx, agent)
  403. const staged = events(agent).filter(event =>
  404. event.type === 'user/message')
  405. expect(staged.map(event => event.type)).toEqual([
  406. 'user/message',
  407. 'user/message',
  408. 'user/message',
  409. ])
  410. expect(JSON.stringify(adapter.requests[0]?.messages)).not.toContain('blocked prompt')
  411. expect(JSON.stringify(adapter.requests[0]?.messages)).toContain('staged context')
  412. expect(JSON.stringify(adapter.requests[0]?.messages)).toContain('staged steering')
  413. })
  414. it('preserves later queued work when a step is rejected', async () => {
  415. const adapter = new MockAdapter([
  416. textResponse('continued'),
  417. textResponse('wake reply'),
  418. ])
  419. const ctx = await harness(adapter)
  420. const agent = await ctx.agentLoop.create(SessionId('rejected-pre-step-order'), {
  421. provider: 'mock',
  422. model: 'mock',
  423. })
  424. ctx.on('agent/pre-step', async ({ messages }, next) => {
  425. const decision = await next()
  426. return messages.some(message =>
  427. message.content.some(block => block.type === 'text' && block.text === 'blocked prompt'))
  428. ? { kind: 'reject' as const }
  429. : decision
  430. })
  431. ctx.on('agent/pre-step', async ({ agent: subject, messages }, next) => {
  432. if (messages.some(message =>
  433. message.content.some(block => block.type === 'text' && block.text === 'blocked prompt'))) {
  434. subject.inject(createUserMessage({
  435. content: [{ type: 'text', text: 'earlier state change' }],
  436. source: { kind: 'plugin', plugin: 'test' },
  437. }))
  438. subject.steer(createUserMessage({
  439. content: [{ type: 'text', text: 'earlier steering' }],
  440. source: { kind: 'user' },
  441. }))
  442. }
  443. return next()
  444. })
  445. const idle = waitForIdle(ctx, agent)
  446. send(agent, 'blocked prompt')
  447. send(agent, 'later prompt')
  448. await idle
  449. expect(events(agent).filter(event => event.type === 'turn/start' || event.type === 'turn/end')
  450. .map(event => event.type)).toEqual(['turn/start', 'turn/end'])
  451. expect(agent.inbox.nextStep.map(message => message.content[0]))
  452. .toEqual([
  453. { type: 'text', text: 'earlier state change' },
  454. { type: 'text', text: 'earlier steering' },
  455. ])
  456. expect(agent.inbox.nextTurn.map(message => message.content[0]))
  457. .toEqual([{ type: 'text', text: 'later prompt' }])
  458. expect(adapter.requests).toEqual([])
  459. const resumed = waitForIdle(ctx, agent)
  460. send(agent, 'wake')
  461. await resumed
  462. const request = JSON.stringify(adapter.requests[0]?.messages)
  463. expect(request).toContain('earlier state change')
  464. expect(request).toContain('earlier steering')
  465. expect(request).toContain('later prompt')
  466. expect(request).not.toContain('blocked prompt')
  467. })
  468. it('preserves context-only injection staged after pre-step began', async () => {
  469. const adapter = new MockAdapter([textResponse('continued')])
  470. const ctx = await harness(adapter)
  471. const agent = await ctx.agentLoop.create(SessionId('rejected-pre-step-context'), { provider: 'mock', model: 'mock' })
  472. const entered = Promise.withResolvers<undefined>()
  473. const decision = Promise.withResolvers<PreStepDecision>()
  474. const disposeBlock = ctx.on('agent/pre-step', async () => {
  475. entered.resolve(undefined)
  476. return decision.promise
  477. })
  478. const idle = waitForIdle(ctx, agent)
  479. send(agent, 'blocked prompt')
  480. await entered.promise
  481. agent.inject(createUserMessage({
  482. content: [{ type: 'text', text: 'independent context' }],
  483. source: { kind: 'plugin', plugin: 'test' },
  484. }))
  485. decision.resolve({ kind: 'reject' })
  486. await idle
  487. const log = events(agent)
  488. expect(log.some(event => event.type === 'user/message')).toBe(false)
  489. expect(agent.inbox.nextStep.map(message => message.content[0]))
  490. .toEqual([{ type: 'text', text: 'independent context' }])
  491. expect(adapter.requests).toEqual([])
  492. disposeBlock()
  493. const resumed = waitForIdle(ctx, agent)
  494. send(agent, 'wake')
  495. await resumed
  496. expect(JSON.stringify(adapter.requests[0]?.messages)).toContain('independent context')
  497. expect(JSON.stringify(adapter.requests[0]?.messages)).not.toContain('blocked prompt')
  498. })
  499. it('leaves inbox state unchanged when its durable append fails', async () => {
  500. const adapter = new MockAdapter([])
  501. const ctx = await harness(adapter)
  502. const agent = await ctx.agentLoop.create(SessionId('rejected-pre-step-append-failure'), {
  503. provider: 'mock',
  504. model: 'mock',
  505. })
  506. vi.spyOn(agent.session, 'append').mockImplementationOnce(() => {
  507. throw new Error('append unavailable')
  508. })
  509. expect(() => {
  510. send(agent, 'blocked prompt')
  511. }).toThrow('append unavailable')
  512. expect(events(agent)).toEqual([])
  513. expect(agent.inbox.hasPending).toBe(false)
  514. expect(agent.status).toBe('idle')
  515. })
  516. it('a blocked prompt preserves adjacent queued prompts', async () => {
  517. const adapter = new MockAdapter([
  518. textResponse('safe reply'),
  519. textResponse('wake reply'),
  520. ])
  521. const ctx = await harness(adapter)
  522. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  523. ctx.on('agent/pre-step', async ({ messages }, next): Promise<PreStepDecision> => {
  524. const text = messages.flatMap(message => message.content)
  525. .map(b => (b.type === 'text' ? b.text : '')).join('')
  526. return text === 'secret'
  527. ? { kind: 'reject' }
  528. : next()
  529. })
  530. const reasons: TurnEndReason[] = []
  531. ctx.on('session/event', (_s, event: SessionEvent) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  532. send(agent, 'secret')
  533. send(agent, 'safe')
  534. await waitForIdle(ctx, agent)
  535. const log = events(agent)
  536. expect(log.filter(e => e.type === 'user/message')).toHaveLength(0)
  537. expect(adapter.requests).toHaveLength(0)
  538. expect(log.filter(e => e.type === 'turn/start')).toHaveLength(1)
  539. expect(log.filter(e => e.type === 'turn/end')).toHaveLength(1)
  540. expect(reasons).toEqual([{ kind: 'blocked' }])
  541. expect(agent.inbox.nextTurn.map(message => message.content[0]))
  542. .toEqual([{ type: 'text', text: 'safe' }])
  543. const resumed = waitForIdle(ctx, agent)
  544. send(agent, 'wake')
  545. await resumed
  546. expect(JSON.stringify(adapter.requests[0]?.messages)).toContain('safe')
  547. expect(JSON.stringify(adapter.requests[0]?.messages)).not.toContain('secret')
  548. })
  549. it('a throwing pre-step listener reports the driver error and retains adjacent work', async () => {
  550. const adapter = new MockAdapter([textResponse('after')])
  551. const ctx = await harness(adapter)
  552. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  553. let threw = false
  554. ctx.on('agent/pre-step', async ({ messages }) => {
  555. if (!threw) { threw = true; throw new Error('prompt hook broke') }
  556. return { kind: 'enter' as const, messages }
  557. })
  558. const errors: Error[] = []
  559. const reasons: TurnEndReason[] = []
  560. const statuses: string[] = []
  561. ctx.on('agent/error', ({ error }) => {
  562. if (error instanceof Error) errors.push(error)
  563. })
  564. ctx.on('agent/status', ({ agent: subject, status }) => { if (subject === agent) statuses.push(status) })
  565. ctx.on('session/event', (session, event) => {
  566. if (session === agent.session && event.type === 'turn/end') reasons.push(event.data.reason)
  567. })
  568. const idle = waitForIdle(ctx, agent)
  569. send(agent, 'first')
  570. send(agent, 'second')
  571. await idle
  572. expect(errors).toEqual([expect.objectContaining({ message: 'prompt hook broke' })])
  573. const log = events(agent)
  574. expect(log.filter(e => e.type === 'turn/start')).toHaveLength(1)
  575. expect(log.filter(e => e.type === 'turn/end')).toHaveLength(1)
  576. expect(reasons).toEqual([{
  577. kind: 'error',
  578. error: { message: 'prompt hook broke', code: 'UNKNOWN' },
  579. }])
  580. expect(statuses).toEqual(['running', 'idle'])
  581. expect(adapter.requests).toHaveLength(0)
  582. expect(agent.inbox.nextTurn.map(message => message.content[0]))
  583. .toEqual([{ type: 'text', text: 'second' }])
  584. })
  585. })
  586. describe('agent/session-start', () => {
  587. it('fires once with source "startup" for a fresh create, before the first turn', async () => {
  588. const adapter = new MockAdapter([textResponse('ok')])
  589. const ctx = await harness(adapter)
  590. const sources: SessionStartSource[] = []
  591. ctx.on('agent/session-start', ({ source }) => void sources.push(source))
  592. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  593. // fires synchronously at create, before any turn
  594. expect(sources).toEqual(['startup'])
  595. expect(events(agent).some(e => e.type === 'turn/start')).toBe(false)
  596. send(agent, 'go')
  597. await waitForIdle(ctx, agent)
  598. // still only one session-start
  599. expect(sources).toEqual(['startup'])
  600. })
  601. it('a session-start listener can inject context the first request sees', async () => {
  602. const adapter = new MockAdapter([textResponse('ok')])
  603. const ctx = await harness(adapter)
  604. ctx.on('agent/session-start', ({ agent }) => {
  605. agent.inject(createUserMessage({ content: [{ type: 'text', text: 'session preamble' }], source: { kind: 'plugin', plugin: 'test' } }))
  606. })
  607. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  608. send(agent, 'go')
  609. await waitForIdle(ctx, agent)
  610. // the injected context reached the model on the first (only) request
  611. expect(JSON.stringify(adapter.requests[0]!.messages)).toContain('session preamble')
  612. // and is recorded with the plugin source, never mislabeled as a user prompt
  613. const ctxMsg = events(agent).find(e => e.type === 'user/message' && e.data.source.kind === 'plugin')
  614. expect(ctxMsg?.type === 'user/message' && ctxMsg.data.source).toEqual({ kind: 'plugin', plugin: 'test' })
  615. })
  616. it('a throwing session-start listener does not abort agent construction', async () => {
  617. const adapter = new MockAdapter([textResponse('ok')])
  618. const ctx = await harness(adapter)
  619. ctx.on('agent/session-start', () => { throw new Error('session-start hook broke') })
  620. // create must not throw — the listener error is contained/logged
  621. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  622. expect(agent.id).toBe(SessionId('a1'))
  623. // and the agent still runs
  624. send(agent, 'go')
  625. await waitForIdle(ctx, agent)
  626. expect(adapter.requests).toHaveLength(1)
  627. })
  628. })
  629. describe('tool additionalContexts buffering across a step', () => {
  630. it('appends each call\'s contexts only AFTER all tool/results, preserving adjacency', async () => {
  631. // One assistant step with TWO tool calls; the second model response stops.
  632. const twoCalls = [
  633. { type: 'block-start' as const, index: 0, blockType: 'tool-call' as const },
  634. { type: 'block-end' as const, index: 0, block: { type: 'tool-call' as const, id: ToolCallId('c1'), name: 'echo', arguments: '{"text":"a"}' } },
  635. { type: 'block-start' as const, index: 1, blockType: 'tool-call' as const },
  636. { type: 'block-end' as const, index: 1, block: { type: 'tool-call' as const, id: ToolCallId('c2'), name: 'echo', arguments: '{"text":"b"}' } },
  637. { type: 'usage' as const, usage: { inputTokens: 5, outputTokens: 5 } },
  638. { type: 'finish' as const, reason: { kind: 'tool-calls' as const } },
  639. ]
  640. const adapter = new MockAdapter([twoCalls, textResponse('done')])
  641. const ctx = await harness(adapter)
  642. ctx.tools.register(defineContentToolFixture({
  643. name: 'echo', description: 'echo', parameters: { text: { type: 'string' } },
  644. async execute(args) { return [{ type: 'text', text: String(args.text) }] },
  645. }))
  646. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  647. // Each call attaches one context naming itself.
  648. ctx.on('tools/post-execute', async (exec, _result): Promise<PostToolDecision> =>
  649. ({
  650. kind: 'accept',
  651. additionalContexts: [createUserMessage({
  652. content: [{ type: 'text', text: `ctx-${exec.callId}` }],
  653. source: { kind: 'plugin', plugin: 'p' },
  654. })],
  655. }))
  656. send(agent, 'go')
  657. await waitForIdle(ctx, agent)
  658. // Event order in the log: both tool/results, THEN both injected contexts —
  659. // never interleaved (which would break tool-call/result adjacency).
  660. const injected = events(agent).filter(e => e.type === 'user/message' && e.data.source.kind === 'plugin')
  661. const seqs = events(agent)
  662. const firstResult = seqs.findIndex(e => e.type === 'tool/result')
  663. const lastResult = seqs.map(e => e.type).lastIndexOf('tool/result')
  664. const firstCtx = seqs.findIndex(e => e === injected[0])
  665. expect(firstResult).toBeGreaterThanOrEqual(0)
  666. expect(lastResult).toBeGreaterThan(firstResult) // two results
  667. expect(firstCtx).toBeGreaterThan(lastResult) // context only after ALL results
  668. // both contexts present
  669. const ctxTexts = injected
  670. .flatMap(e => (e.type === 'user/message' ? e.data.content : []))
  671. .map(b => (b.type === 'text' ? b.text : ''))
  672. expect(ctxTexts).toEqual(['ctx-c1', 'ctx-c2'])
  673. })
  674. it('appends multiple contexts deferred by one composite tool after its outer result', async () => {
  675. const adapter = new MockAdapter([toolCallResponse('c1', 'composite', {}), textResponse('done')])
  676. const ctx = await harness(adapter)
  677. ctx.tools.register(defineContentToolFixture({
  678. name: 'composite', description: 'composite', parameters: {},
  679. async execute(_args, exec) {
  680. exec.deferContext(createUserMessage({
  681. content: [{ type: 'text', text: 'nested-a' }], source: { kind: 'plugin', plugin: 'a' },
  682. }))
  683. exec.deferContext(createUserMessage({
  684. content: [{ type: 'text', text: 'nested-b' }], source: { kind: 'plugin', plugin: 'b' },
  685. }))
  686. return [{ type: 'text', text: 'outer result' }]
  687. },
  688. }))
  689. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  690. send(agent, 'go')
  691. await waitForIdle(ctx, agent)
  692. const log = events(agent)
  693. const resultIndex = log.findIndex(event => event.type === 'tool/result')
  694. const contextEvents = log.filter(event => event.type === 'user/message' && event.data.source.kind === 'plugin')
  695. expect(resultIndex).toBeGreaterThanOrEqual(0)
  696. expect(log.findIndex(event => event === contextEvents[0])).toBeGreaterThan(resultIndex)
  697. expect(contextEvents.map(event => event.type === 'user/message' && event.data.source)).toEqual([
  698. { kind: 'plugin', plugin: 'a' },
  699. { kind: 'plugin', plugin: 'b' },
  700. ])
  701. })
  702. })
  703. describe('tools/pre-execute gate (native-plugin permission pattern, end-to-end through the loop)', () => {
  704. it('deny short-circuits dispatch into an isError result the model sees', async () => {
  705. const adapter = new MockAdapter([toolCallResponse('c1', 'danger', {}), textResponse('ok')])
  706. const ctx = await harness(adapter)
  707. let ran = false
  708. ctx.tools.register(defineContentToolFixture({
  709. name: 'danger', description: 'danger', parameters: {},
  710. async execute() { ran = true; return [{ type: 'text', text: 'should not run' }] },
  711. }))
  712. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  713. ctx.on('tools/pre-execute', async (exec, next): Promise<PreToolDecision> => {
  714. if (exec.name === 'danger') return { kind: 'deny', reason: 'blocked dangerous tool' }
  715. return next()
  716. })
  717. send(agent, 'go')
  718. await waitForIdle(ctx, agent)
  719. expect(ran).toBe(false)
  720. const result = events(agent).find(e => e.type === 'tool/result')
  721. expect(result?.type === 'tool/result' && result.data.message.content[0].isError).toBe(true)
  722. expect(result?.type === 'tool/result'
  723. && result.data.message.content[0].content.some(b => b.type === 'text' && b.text.includes('blocked dangerous tool'))).toBe(true)
  724. })
  725. })
  726. describe('worked example: a native hook plugin is just a cordis plugin on the seams', () => {
  727. // The whole point of the interception taxonomy: a "native hook" needs no dsh-hook-protocol,
  728. // no external command, no hook/* log — it is an ordinary cordis plugin subscribing to the
  729. // canonical events and returning typed decisions.
  730. const NativeGuard = {
  731. name: 'native-guard',
  732. apply(ctx: Context) {
  733. // 1. SessionStart: seed a standing instruction.
  734. ctx.on('agent/session-start', ({ agent, source }) => {
  735. agent.inject(createUserMessage({ content: [{ type: 'text', text: `policy active (started: ${source})` }], source: { kind: 'plugin', plugin: 'native-guard' } }))
  736. })
  737. // 2. PreStep: reject a forbidden prompt, annotate the rest.
  738. ctx.on('agent/pre-step', async ({ messages }, next): Promise<PreStepDecision> => {
  739. const text = messages.flatMap(message => message.content)
  740. .map(b => (b.type === 'text' ? b.text : '')).join('')
  741. if (text.includes('rm -rf')) {
  742. return { kind: 'reject' }
  743. }
  744. return next()
  745. })
  746. // 3. PreToolUse: deny a dangerous tool by name.
  747. ctx.on('tools/pre-execute', async (exec, next): Promise<PreToolDecision> => {
  748. if (exec.name === 'danger') return { kind: 'deny', reason: 'danger tool denied' }
  749. return next()
  750. })
  751. // 4. PostToolUse: attach context after a tool runs.
  752. ctx.on('tools/post-execute', async (_exec, _result, next): Promise<PostToolDecision> => {
  753. const decision = await next()
  754. if (decision.kind === 'accept') {
  755. return { kind: 'accept', additionalContexts: [createUserMessage({
  756. content: [{ type: 'text', text: 'audited' }], source: { kind: 'plugin', plugin: 'native-guard' },
  757. })] }
  758. }
  759. return decision
  760. })
  761. },
  762. }
  763. it('all four seams fire for a real allowed turn with a tool call', async () => {
  764. const adapter = new MockAdapter([toolCallResponse('c1', 'echo', { text: 'hi' }), textResponse('done')])
  765. const ctx = await harness(adapter)
  766. await ctx.plugin(NativeGuard)
  767. ctx.tools.register(defineContentToolFixture({
  768. name: 'echo', description: 'echo', parameters: { text: { type: 'string' } },
  769. async execute(args) { return [{ type: 'text', text: String(args.text) }] },
  770. }))
  771. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  772. send(agent, 'please echo hi')
  773. await waitForIdle(ctx, agent)
  774. const log = events(agent)
  775. // session-start preamble injected
  776. expect(log.some(e => e.type === 'user/message' && e.data.source.kind === 'plugin'
  777. && e.data.content.some(b => b.type === 'text' && b.text.includes('policy active (started: startup)')))).toBe(true)
  778. // prompt allowed → user-sourced user/message recorded
  779. expect(log.some(e => e.type === 'user/message' && e.data.source.kind === 'user')).toBe(true)
  780. // tool ran (echo allowed) and post-execute attached "audited" context
  781. expect(log.some(e => e.type === 'tool/result' && !e.data.message.content[0].isError)).toBe(true)
  782. expect(log.some(e => e.type === 'user/message' && e.data.source.kind === 'plugin'
  783. && e.data.content.some(b => b.type === 'text' && b.text === 'audited'))).toBe(true)
  784. // NO hook/* events — a native plugin needs none
  785. expect(log.some(e => e.type.startsWith('hook/'))).toBe(false)
  786. })
  787. it('the same plugin blocks a destructive prompt inside a no-step turn', async () => {
  788. const adapter = new MockAdapter([textResponse('should not run')])
  789. const ctx = await harness(adapter)
  790. await ctx.plugin(NativeGuard)
  791. const agent = await ctx.agentLoop.create(SessionId('a2'), { provider: 'mock', model: 'mock' })
  792. const reasons: TurnEndReason[] = []
  793. ctx.on('session/event', (_s, event: SessionEvent) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  794. send(agent, 'run rm -rf /')
  795. await agent.whenIdle()
  796. expect(adapter.requests).toHaveLength(0)
  797. expect(reasons).toEqual([{ kind: 'blocked' }])
  798. })
  799. it('HMR-safety: disposing the plugin fiber removes all four listeners', async () => {
  800. const adapter = new MockAdapter([textResponse('ok')])
  801. const ctx = await harness(adapter)
  802. const fiber = await ctx.plugin(NativeGuard)
  803. await fiber.dispose()
  804. // After disposal, a destructive prompt is NOT blocked (the listener is gone).
  805. const agent = await ctx.agentLoop.create(SessionId('a3'), { provider: 'mock', model: 'mock' })
  806. send(agent, 'run rm -rf /')
  807. await waitForIdle(ctx, agent)
  808. // the prompt ran (not rejected) — proving the pre-step listener was disposed
  809. expect(adapter.requests).toHaveLength(1)
  810. expect(events(agent).some(e => e.type === 'user/message')).toBe(true)
  811. })
  812. })