api-proxy-commands.spec.ts 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579
  1. import { MessageId, freezeMessage } from '@deepseek-ai/dsh-llm'
  2. /**
  3. * Command/skill RPC handlers and the two new frames over createApiProxy:
  4. * command.list serves the addressed agent's effective catalog (missing
  5. * registry = loud internal error), command.execute dispatches through the
  6. * registry with the carrier signal, skill.list resolves cwd from the session
  7. * header (never via the Agent registry), the host stream broadcasts
  8. * commands-changed, and the mux stream carries live queued frames plus the
  9. * open-time queue snapshot.
  10. */
  11. import { describe, expect, it, vi } from 'vitest'
  12. import { Context } from 'cordis'
  13. import AgentRegistry, { InboxItemId } from '@deepseek-ai/dsh-agent'
  14. import type { Agent, InboxItem, InboxPlacement } from '@deepseek-ai/dsh-agent'
  15. import SessionStore from '@deepseek-ai/dsh-session'
  16. import type { SessionId, UserMessage } from '@deepseek-ai/dsh-session'
  17. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  18. import ToolRegistry from '@deepseek-ai/dsh-tools'
  19. import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
  20. import CommandService from '@deepseek-ai/dsh-commands'
  21. import SkillService from '@deepseek-ai/dsh-skill'
  22. import type { HostFrame, MuxFrame } from '../src/api/index.ts'
  23. import type { RpcRequest, RpcResponse } from '../src/api/rpc.ts'
  24. import { RpcId } from '../src/api/rpc.ts'
  25. import { createApiProxy } from '../src/api-proxy.ts'
  26. const DEFAULTS = { provider: 'p', model: 'm', cwd: '/tmp', workspaceRoot: '/tmp' }
  27. function request<P>(payload: P): RpcRequest<P> {
  28. return { rpcId: RpcId(`req-${String(nextRpc++)}`), payload }
  29. }
  30. let nextRpc = 1
  31. function expectOk<T>(response: RpcResponse<T>): T {
  32. expect(response.result.ok).toBe(true)
  33. if (!response.result.ok) throw new Error('unreachable')
  34. return response.result.value
  35. }
  36. function expectErr<T>(response: RpcResponse<T>): { code: string; message: string } {
  37. expect(response.result.ok).toBe(false)
  38. if (response.result.ok) throw new Error('unreachable')
  39. return response.result.error
  40. }
  41. /** Composition floor for the command/skill paths (no LLM, no persistence). */
  42. async function harness(options: { commands?: boolean; skills?: boolean } = {}): Promise<Context> {
  43. const ctx = new Context()
  44. await ctx.plugin(SessionStore)
  45. await ctx.plugin(SystemPrompt, { persona: '' })
  46. await ctx.plugin(ToolRegistry)
  47. await ctx.plugin(UserInteractionService)
  48. await ctx.plugin(AgentRegistry)
  49. if (options.skills !== false) await ctx.plugin(SkillService, {})
  50. if (options.commands !== false) await ctx.plugin(CommandService)
  51. // Host-stream opener reads the committed-workspace baseline; the stub
  52. // suffices here — the real workspace composition is api-proxy-workspace.spec's.
  53. ctx.provide('workspace', { list: () => [] } as never)
  54. return ctx
  55. }
  56. /** Register a live structural agent stub (api-proxy-view precedent: only id/session/status/ctx are read). */
  57. function stubAgent(ctx: Context, sessionId?: SessionId): Agent {
  58. const session = ctx.sessions.create(sessionId)
  59. const agent = { id: session.id, session, status: 'idle', ctx } as Agent
  60. ctx.agents.register(agent)
  61. return agent
  62. }
  63. /** Drain `count` frames from a stream, then abort it. */
  64. async function collect<F>(iterable: AsyncIterable<RpcRequest<F>>, count: number, abort: AbortController): Promise<F[]> {
  65. const frames: F[] = []
  66. for await (const frame of iterable) {
  67. frames.push(frame.payload)
  68. if (frames.length >= count) abort.abort()
  69. }
  70. return frames
  71. }
  72. describe('command.list', () => {
  73. it('serves the addressed agent\'s name-sorted catalog', async () => {
  74. const ctx = await harness()
  75. ctx.commands.register({ name: 'zeta', description: 'z', handler: () => ({ kind: 'success' }) })
  76. ctx.commands.register({ name: 'alpha', description: 'a', input: { hint: '<x>' }, handler: () => ({ kind: 'success' }) })
  77. const api = createApiProxy(ctx, DEFAULTS)
  78. const agent = stubAgent(ctx)
  79. const value = expectOk(await api.commands.list(request({ sessionId: agent.id })))
  80. expect(value.commands).toEqual([
  81. { name: 'alpha', description: 'a', input: { hint: '<x>' } },
  82. { name: 'zeta', description: 'z' },
  83. ])
  84. })
  85. it('fails loud with internal when the command registry is not mounted', async () => {
  86. const ctx = await harness({ commands: false })
  87. const api = createApiProxy(ctx, DEFAULTS)
  88. const error = expectErr(await api.commands.list(request({ sessionId: 's' as SessionId })))
  89. expect(error.code).toBe('internal')
  90. expect(error.message).toContain('command registry')
  91. })
  92. it('does not route a live subagent through the generic command domain', async () => {
  93. const ctx = await harness()
  94. const session = ctx.sessions.create(undefined, { meta: { cwd: '/proj', origin: 'subagent' } })
  95. const agent = { id: session.id, session, status: 'idle', ctx } as Agent
  96. ctx.agents.register(agent)
  97. const api = createApiProxy(ctx, DEFAULTS)
  98. const error = expectErr(await api.commands.list(request({ sessionId: agent.id })))
  99. expect(error).toMatchObject({ code: 'agent-busy' })
  100. })
  101. })
  102. describe('command.execute', () => {
  103. it('executes a known command against the addressed agent and detaches the result', async () => {
  104. const ctx = await harness()
  105. let received: string | undefined
  106. ctx.commands.register({
  107. name: 'goal',
  108. description: 'set goal',
  109. handler: (invocation) => {
  110. received = invocation.rawInput
  111. return { kind: 'success', text: `goal:${invocation.agent.id}` }
  112. },
  113. })
  114. const api = createApiProxy(ctx, DEFAULTS)
  115. const agent = stubAgent(ctx)
  116. const value = expectOk(await api.commands.execute(request({ sessionId: agent.id, line: '/goal ship it' }), new AbortController().signal))
  117. expect(value).toMatchObject({ matched: true })
  118. expect(value.commandId).toBeTruthy()
  119. expect(received).toBe(' ship it')
  120. // Pure admission on the wire: the outcome rides the durably logged
  121. // lifecycle pair instead of the response.
  122. const lifecycle = agent.session.events.filter(e => e.type === 'command/run' || e.type === 'command/done')
  123. expect(lifecycle).toMatchObject([
  124. { type: 'command/run', data: { commandId: value.commandId, name: 'goal', args: ' ship it' } },
  125. { type: 'command/done', data: { commandId: value.commandId, kind: 'success', text: `goal:${agent.id}` } },
  126. ])
  127. })
  128. it('returns matched:false when syntax or name does not resolve', async () => {
  129. const ctx = await harness()
  130. const api = createApiProxy(ctx, DEFAULTS)
  131. const agent = stubAgent(ctx)
  132. const signal = new AbortController().signal
  133. expect(expectOk(await api.commands.execute(request({ sessionId: agent.id, line: '/unknown' }), signal))).toEqual({ matched: false })
  134. expect(expectOk(await api.commands.execute(request({ sessionId: agent.id, line: 'not a command' }), signal))).toEqual({ matched: false })
  135. })
  136. it('maps a session miss to session-not-found and a registry gap to internal', async () => {
  137. const ctx = await harness()
  138. const api = createApiProxy(ctx, DEFAULTS)
  139. const missing = expectErr(await api.commands.execute(
  140. request({ sessionId: 'session-nope' as SessionId, line: '/x' }), new AbortController().signal))
  141. expect(missing.code).toBe('internal') // Cold Agent-bound access fails loud when persistence is absent.
  142. const bare = await harness({ commands: false })
  143. const bareApi = createApiProxy(bare, DEFAULTS)
  144. expect(expectErr(await bareApi.commands.execute(
  145. request({ sessionId: 's' as SessionId, line: '/x' }), new AbortController().signal)).code).toBe('internal')
  146. })
  147. it('reports an aborted handler as cancelled and a throwing handler as internal', async () => {
  148. const ctx = await harness()
  149. ctx.commands.register({
  150. name: 'hang',
  151. description: 'never settles on its own',
  152. handler: () => new Promise(() => { /* settled only by abort */ }),
  153. })
  154. ctx.commands.register({
  155. name: 'boom',
  156. description: 'throws',
  157. handler: () => { throw new Error('kaboom') },
  158. })
  159. const api = createApiProxy(ctx, DEFAULTS)
  160. const agent = stubAgent(ctx)
  161. const controller = new AbortController()
  162. const pending = api.commands.execute(request({ sessionId: agent.id, line: '/hang' }), controller.signal)
  163. controller.abort()
  164. expect(expectErr(await pending).code).toBe('cancelled')
  165. const thrown = expectErr(await api.commands.execute(request({ sessionId: agent.id, line: '/boom' }), new AbortController().signal))
  166. expect(thrown.code).toBe('internal')
  167. expect(thrown.message).toContain('kaboom')
  168. })
  169. })
  170. describe('skill.list', () => {
  171. it('lists skills for the session cwd taken from the header', async () => {
  172. const ctx = await harness()
  173. const seenCwds: (string | undefined)[] = []
  174. ctx.skills.registerProvider(() => ({
  175. name: 'probe',
  176. list: (options) => {
  177. seenCwds.push(options.cwd)
  178. return Promise.resolve([
  179. {
  180. name: 'commit-helper', description: 'Git commits', whenToUse: 'when committing',
  181. invocation: { modelInvocable: true, userInvocable: true },
  182. source: 'custom', provider: 'probe', rank: 0, locator: null,
  183. },
  184. {
  185. name: 'user-only', description: 'User-only',
  186. invocation: { modelInvocable: false, userInvocable: true },
  187. source: 'custom', provider: 'probe', rank: 0, locator: null,
  188. },
  189. {
  190. name: 'model-only', description: 'Model-only',
  191. invocation: { modelInvocable: true, userInvocable: false },
  192. source: 'custom', provider: 'probe', rank: 0, locator: null,
  193. },
  194. {
  195. name: 'trusted-only', description: 'Trusted-only',
  196. invocation: { modelInvocable: false, userInvocable: false },
  197. source: 'custom', provider: 'probe', rank: 0, locator: null,
  198. },
  199. ])
  200. },
  201. get: () => Promise.resolve(undefined),
  202. }))
  203. const api = createApiProxy(ctx, DEFAULTS)
  204. // No agent is registered for this session: header resolution must not
  205. // touch (or resume through) the Agent registry.
  206. const session = ctx.sessions.create(undefined, { meta: { cwd: '/proj' } })
  207. const value = expectOk(await api.skills.list(request({ sessionId: session.id })))
  208. expect(value.skills).toEqual([{ name: 'commit-helper', description: 'Git commits', whenToUse: 'when committing' }])
  209. expect(seenCwds).toEqual(['/proj'])
  210. expect(ctx.agents.get(session.id)).toBeUndefined()
  211. })
  212. it('fails loud on an unattached session id (business error, no resume attempt)', async () => {
  213. const ctx = await harness()
  214. const api = createApiProxy(ctx, DEFAULTS)
  215. const error = expectErr(await api.skills.list(request({ sessionId: 'session-cold' as SessionId })))
  216. expect(error.code).toBe('session-not-found')
  217. })
  218. it('fails loud with internal when the skill registry is not mounted', async () => {
  219. const ctx = await harness({ skills: false })
  220. const api = createApiProxy(ctx, DEFAULTS)
  221. const session = ctx.sessions.create(undefined, { meta: { cwd: '/proj' } })
  222. const error = expectErr(await api.skills.list(request({ sessionId: session.id })))
  223. expect(error.code).toBe('internal')
  224. expect(error.message).toContain('skill registry is absent')
  225. })
  226. it('folds a provider failure into internal', async () => {
  227. const ctx = await harness()
  228. ctx.skills.registerProvider(() => ({
  229. name: 'broken',
  230. list: () => Promise.reject(new Error('directory exploded')),
  231. get: () => Promise.resolve(undefined),
  232. }))
  233. const api = createApiProxy(ctx, DEFAULTS)
  234. const session = ctx.sessions.create(undefined, { meta: { cwd: '/proj' } })
  235. const response = await api.skills.list(request({ sessionId: session.id }))
  236. // dsh-skill contains one provider's failure (logs and serves the rest), so
  237. // this surfaces as an empty ok catalog rather than an error.
  238. const value = expectOk(response)
  239. expect(value.skills).toEqual([])
  240. })
  241. })
  242. describe('host/commands-changed frame', () => {
  243. it('broadcasts on registry change', async () => {
  244. const ctx = await harness()
  245. const api = createApiProxy(ctx, DEFAULTS)
  246. const abort = new AbortController()
  247. const stream = api.events.host({ rpcId: RpcId('t-host'), payload: {} }, abort.signal)
  248. const collected = collect<HostFrame>(stream, 1, abort)
  249. ctx.commands.register({ name: 'late', description: 'l', handler: () => ({ kind: 'success' }) })
  250. expect(await collected).toEqual([{ type: 'host/commands-changed' }])
  251. })
  252. })
  253. /** Build one frozen inbox message for the live `agent/inbox/*` events. */
  254. function inboxMessage(id: string, text: string, rpcId?: string): UserMessage {
  255. return freezeMessage({
  256. id: MessageId(id),
  257. role: 'user',
  258. content: [{ type: 'text' as const, text }],
  259. source: rpcId === undefined ? { kind: 'user' as const } : { kind: 'user' as const, rpcId: RpcId(rpcId) },
  260. })
  261. }
  262. /** Build one addressable inbox occurrence around a frozen message. */
  263. function inboxItem(id: string, message: UserMessage, placement: InboxPlacement): InboxItem {
  264. return { id: InboxItemId(id), message, placement }
  265. }
  266. describe('session.updateQueue', () => {
  267. it('routes addressable actions and reports strict steer races', async () => {
  268. const ctx = await harness()
  269. const agent = stubAgent(ctx)
  270. const seen: unknown[] = []
  271. agent.updateInbox = (id, action) => {
  272. seen.push({ id, action })
  273. if (id === InboxItemId('present')) return 'applied'
  274. return id === InboxItemId('closed') ? 'steer-unavailable' : 'not-found'
  275. }
  276. const api = createApiProxy(ctx, DEFAULTS)
  277. const applied = await api.sessions.updateQueue({
  278. rpcId: RpcId('q-apply'),
  279. payload: {
  280. sessionId: agent.id,
  281. itemId: InboxItemId('present'),
  282. action: { kind: 'edit', content: [{ type: 'text', text: 'edited' }] },
  283. },
  284. })
  285. expect(expectOk(applied)).toEqual({ accepted: true })
  286. const missing = await api.sessions.updateQueue({
  287. rpcId: RpcId('q-missing'),
  288. payload: {
  289. sessionId: agent.id,
  290. itemId: InboxItemId('claimed'),
  291. action: { kind: 'remove' },
  292. },
  293. })
  294. expect(expectErr(missing)).toMatchObject({ code: 'queue-item-not-found' })
  295. const closed = await api.sessions.updateQueue({
  296. rpcId: RpcId('q-closed'),
  297. payload: {
  298. sessionId: agent.id,
  299. itemId: InboxItemId('closed'),
  300. action: { kind: 'steer' },
  301. },
  302. })
  303. expect(expectErr(closed)).toMatchObject({
  304. code: 'steer-unavailable',
  305. details: { itemId: 'closed' },
  306. })
  307. expect(seen).toEqual([
  308. { id: 'present', action: { kind: 'edit', content: [{ type: 'text', text: 'edited' }] } },
  309. { id: 'claimed', action: { kind: 'remove' } },
  310. { id: 'closed', action: { kind: 'steer' } },
  311. ])
  312. })
  313. it('rejects a stale occurrence without resuming a cold agent', async () => {
  314. const ctx = await harness()
  315. const resume = vi.spyOn(ctx.agents, 'resume')
  316. const api = createApiProxy(ctx, DEFAULTS)
  317. const response = await api.sessions.updateQueue({
  318. rpcId: RpcId('q-cold'),
  319. payload: {
  320. sessionId: 'cold-session' as SessionId,
  321. itemId: InboxItemId('stale-item'),
  322. action: { kind: 'remove' },
  323. },
  324. })
  325. expect(expectErr(response)).toMatchObject({ code: 'queue-item-not-found' })
  326. expect(resume).not.toHaveBeenCalled()
  327. })
  328. })
  329. describe('session/queue frames', () => {
  330. it('folds nested mutations observed before their outer enqueue', async () => {
  331. const ctx = await harness()
  332. const agent = stubAgent(ctx)
  333. const original = inboxItem('i-edit', inboxMessage('m-edit', 'before'), 'queued')
  334. const edited = inboxItem('i-edit', inboxMessage('m-edit', 'after'), 'queued')
  335. const removed = inboxItem('i-remove', inboxMessage('m-remove', 'remove me'), 'queued')
  336. ctx.on('agent/inbox/enqueue', (subject, item) => {
  337. if (subject !== agent) return
  338. if (item.id === original.id) ctx.emit('agent/inbox/update', agent, edited)
  339. if (item.id === removed.id) ctx.emit('agent/inbox/discard', agent, [removed])
  340. })
  341. const api = createApiProxy(ctx, DEFAULTS)
  342. const live = new AbortController()
  343. const collected = collect<MuxFrame>(
  344. api.events.mux({ rpcId: RpcId('t-mux-reentrant'), payload: {} }, live.signal), 2, live)
  345. ctx.emit('agent/inbox/enqueue', agent, original)
  346. ctx.emit('agent/inbox/enqueue', agent, removed)
  347. const liveFrames = (await collected).filter(frame => frame.type === 'session/queue')
  348. expect(liveFrames.map(frame => frame.items)).toEqual([
  349. [{ id: edited.id, placement: edited.placement, message: edited.message }],
  350. ])
  351. const replay = new AbortController()
  352. const replayFrames = await collect<MuxFrame>(
  353. api.events.mux({ rpcId: RpcId('t-mux-reentrant-replay'), payload: {} }, replay.signal), 2, replay)
  354. expect(replayFrames.filter(frame => frame.type === 'session/queue')).toEqual(liveFrames)
  355. })
  356. it('expires unmatched mutations after the synchronous re-entry window', async () => {
  357. const ctx = await harness()
  358. const agent = stubAgent(ctx)
  359. const api = createApiProxy(ctx, DEFAULTS)
  360. const original = inboxItem('i-stale-edit', inboxMessage('m-stale-edit', 'original'), 'queued')
  361. const staleEdit = inboxItem('i-stale-edit', inboxMessage('m-stale-edit', 'stale edit'), 'queued')
  362. const staleTerminal = inboxItem('i-stale-terminal', inboxMessage('m-stale-terminal', 'keep me'), 'queued')
  363. ctx.emit('agent/inbox/update', agent, staleEdit)
  364. ctx.emit('agent/inbox/discard', agent, [staleTerminal])
  365. await Promise.resolve()
  366. ctx.emit('agent/inbox/enqueue', agent, original)
  367. ctx.emit('agent/inbox/enqueue', agent, staleTerminal)
  368. const replay = new AbortController()
  369. const frames = await collect<MuxFrame>(
  370. api.events.mux({ rpcId: RpcId('t-mux-expired-unseen'), payload: {} }, replay.signal), 2, replay)
  371. expect(frames.filter(frame => frame.type === 'session/queue')).toEqual([{
  372. type: 'session/queue',
  373. sessionId: agent.id,
  374. items: [
  375. { id: original.id, placement: original.placement, message: original.message },
  376. { id: staleTerminal.id, placement: staleTerminal.placement, message: staleTerminal.message },
  377. ],
  378. }])
  379. })
  380. it('publishes complete live snapshots and replays the latest snapshot on reconnect', async () => {
  381. const ctx = await harness()
  382. const api = createApiProxy(ctx, DEFAULTS)
  383. const agent = stubAgent(ctx)
  384. const live = new AbortController()
  385. const liveStream = api.events.mux({ rpcId: RpcId('t-mux-live'), payload: {} }, live.signal)
  386. // subscribed baseline + one snapshot per accepted inbox occurrence.
  387. const liveCollected = collect<MuxFrame>(liveStream, 3, live)
  388. const queued = inboxItem('i-1', inboxMessage('m-1', 'queued prompt'), 'queued')
  389. const steering = inboxItem('i-2', inboxMessage('m-2', 'steering prompt'), 'steering')
  390. ctx.emit('agent/inbox/enqueue', agent, queued)
  391. ctx.emit('agent/inbox/enqueue', agent, steering)
  392. const liveFrames = (await liveCollected).filter(f => f.type === 'session/queue')
  393. expect(liveFrames).toEqual([
  394. {
  395. type: 'session/queue',
  396. sessionId: agent.id,
  397. items: [{ id: queued.id, placement: 'queued', message: queued.message }],
  398. },
  399. {
  400. type: 'session/queue',
  401. sessionId: agent.id,
  402. items: [
  403. { id: queued.id, placement: 'queued', message: queued.message },
  404. { id: steering.id, placement: 'steering', message: steering.message },
  405. ],
  406. },
  407. ])
  408. // A fresh mux connection replays only the current authoritative snapshot.
  409. const replay = new AbortController()
  410. const replayFrames = await collect<MuxFrame>(
  411. api.events.mux({ rpcId: RpcId('t-mux-replay'), payload: {} }, replay.signal), 2, replay)
  412. expect(replayFrames.filter(f => f.type === 'session/queue')).toEqual([liveFrames[1]])
  413. })
  414. it('publishes the durable steering event before retiring its transient row', async () => {
  415. const ctx = await harness()
  416. const api = createApiProxy(ctx, DEFAULTS)
  417. const agent = stubAgent(ctx)
  418. const abort = new AbortController()
  419. const collected = collect<MuxFrame>(
  420. api.events.mux({ rpcId: RpcId('t-steering-order'), payload: {} }, abort.signal), 5, abort)
  421. const steering = inboxItem('i-steering', inboxMessage('m-steering', 'interrupt now'), 'steering')
  422. agent.session.append('turn/start', {
  423. turn: 1,
  424. trigger: { kind: 'message', source: { kind: 'user' } },
  425. })
  426. ctx.emit('agent/inbox/enqueue', agent, steering)
  427. ctx.emit('agent/inbox/dequeue', agent, steering)
  428. agent.session.append('steering/message', {
  429. turn: 1,
  430. message: steering.message,
  431. }, { surfaceOp: 'append' })
  432. const frames = await collected
  433. expect(frames.map(frame => frame.type)).toEqual([
  434. 'session/subscribed',
  435. 'session/event',
  436. 'session/queue',
  437. 'session/event',
  438. 'session/queue',
  439. ])
  440. expect(frames[2]).toMatchObject({
  441. type: 'session/queue',
  442. items: [{ id: steering.id, placement: 'steering' }],
  443. })
  444. expect(frames[3]).toMatchObject({
  445. type: 'session/event',
  446. event: { type: 'steering/message', data: { message: { id: steering.message.id } } },
  447. })
  448. expect(frames[4]).toMatchObject({ type: 'session/queue', items: [] })
  449. })
  450. it('retains claimed steering in re-entrant snapshots until its durable event', async () => {
  451. const ctx = await harness()
  452. const api = createApiProxy(ctx, DEFAULTS)
  453. const agent = stubAgent(ctx)
  454. const steering = inboxItem('i-steering', inboxMessage('m-steering', 'interrupt now'), 'steering')
  455. const queued = inboxItem('i-reentrant', inboxMessage('m-reentrant', 'later'), 'queued')
  456. ctx.on('agent/inbox/dequeue', (subject, item) => {
  457. if (subject === agent && item.id === steering.id) ctx.emit('agent/inbox/enqueue', agent, queued)
  458. })
  459. const abort = new AbortController()
  460. const collected = collect<MuxFrame>(
  461. api.events.mux({ rpcId: RpcId('t-steering-reentrant-order'), payload: {} }, abort.signal), 6, abort)
  462. agent.session.append('turn/start', {
  463. turn: 1,
  464. trigger: { kind: 'message', source: { kind: 'user' } },
  465. })
  466. ctx.emit('agent/inbox/enqueue', agent, steering)
  467. ctx.emit('agent/inbox/dequeue', agent, steering)
  468. agent.session.append('steering/message', {
  469. turn: 1,
  470. message: steering.message,
  471. }, { surfaceOp: 'append' })
  472. const frames = await collected
  473. expect(frames[3]).toMatchObject({
  474. type: 'session/queue',
  475. items: [
  476. { id: steering.id, placement: 'steering' },
  477. { id: queued.id, placement: 'queued' },
  478. ],
  479. })
  480. expect(frames[4]).toMatchObject({
  481. type: 'session/event',
  482. event: { type: 'steering/message', data: { message: { id: steering.message.id } } },
  483. })
  484. expect(frames[5]).toMatchObject({
  485. type: 'session/queue',
  486. items: [{ id: queued.id, placement: 'queued' }],
  487. })
  488. })
  489. it('publishes edits in place in the authoritative order', async () => {
  490. const ctx = await harness()
  491. const api = createApiProxy(ctx, DEFAULTS)
  492. const agent = stubAgent(ctx)
  493. const abort = new AbortController()
  494. const collected = collect<MuxFrame>(
  495. api.events.mux({ rpcId: RpcId('t-mux-updates'), payload: {} }, abort.signal), 5, abort)
  496. const first = inboxItem('i-a', inboxMessage('m-a', 'a'), 'queued')
  497. const second = inboxItem('i-b', inboxMessage('m-b', 'b'), 'queued')
  498. const edited = inboxItem('i-b', inboxMessage('m-b', 'b edited'), 'queued')
  499. ctx.emit('agent/inbox/enqueue', agent, first)
  500. ctx.emit('agent/inbox/enqueue', agent, second)
  501. ctx.emit('agent/inbox/update', agent, edited)
  502. ctx.emit('agent/inbox/dequeue', agent, edited)
  503. const frames = (await collected).filter(frame => frame.type === 'session/queue')
  504. expect(frames.map(frame => frame.items)).toEqual([
  505. [{ id: first.id, placement: first.placement, message: first.message }],
  506. [
  507. { id: first.id, placement: first.placement, message: first.message },
  508. { id: second.id, placement: second.placement, message: second.message },
  509. ],
  510. [
  511. { id: first.id, placement: first.placement, message: first.message },
  512. { id: edited.id, placement: edited.placement, message: edited.message },
  513. ],
  514. [{ id: first.id, placement: first.placement, message: first.message }],
  515. ])
  516. })
  517. it('publishes an empty snapshot after terminal discard', async () => {
  518. const ctx = await harness()
  519. const api = createApiProxy(ctx, DEFAULTS)
  520. const agent = stubAgent(ctx)
  521. const doomed = inboxItem('i-doomed', inboxMessage('m-5', 'doomed'), 'queued')
  522. ctx.emit('agent/inbox/enqueue', agent, doomed)
  523. ctx.emit('agent/inbox/discard', agent, [doomed])
  524. const abort = new AbortController()
  525. const frames = await collect<MuxFrame>(
  526. api.events.mux({ rpcId: RpcId('t-mux-swept'), payload: {} }, abort.signal), 1, abort)
  527. expect(frames.filter(frame => frame.type === 'session/queue')).toHaveLength(0)
  528. })
  529. })