commands-queue-attachment.host.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281
  1. import { Context } from '@deepseek-ai/cordis'
  2. import AgentRegistry from '@deepseek-ai/dsh-agent'
  3. import type { Agent, Inbox, ModelSelectionRef } from '@deepseek-ai/dsh-agent'
  4. import { AttachmentError, AttachmentId } from '@deepseek-ai/dsh-attachment'
  5. import type { ImageAttachmentRef } from '@deepseek-ai/dsh-attachment'
  6. import { createAssistantMessage, createUserMessage, MessageId } from '@deepseek-ai/dsh-llm'
  7. import SessionStore, { SessionId, SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'
  8. import type { SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  9. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  10. import { describe, expect, it, vi } from 'vitest'
  11. import { ApiSessionAgentController } from '../src/agent.ts'
  12. import { SessionCommandController } from '../src/commands.ts'
  13. import { createInboxStub } from '@deepseek-ai/dsh-agent-loop-testkit'
  14. import { installSessionReadTestServices, testSessionPersistence } from './test-remote.ts'
  15. async function commandHarness(): Promise<{
  16. ctx: Context
  17. controller: SessionCommandController
  18. agent: Agent
  19. inbox: Inbox
  20. steer: ReturnType<typeof vi.fn>
  21. cancel: ReturnType<typeof vi.fn>
  22. }> {
  23. const ctx = new Context()
  24. await ctx.plugin(SessionStore)
  25. await ctx.plugin(SessionProjectionRegistry)
  26. await ctx.plugin(AgentRegistry)
  27. const session = ctx.sessions.create(SessionId('commands-session'), { meta: { cwd: '/workspace' } })
  28. const inbox = createInboxStub()
  29. const steer = vi.fn()
  30. const cancel = vi.fn()
  31. const agent: Agent = {
  32. id: session.id,
  33. options: {},
  34. session,
  35. inbox,
  36. status: 'running',
  37. ctx,
  38. send: () => {},
  39. followup: vi.fn(),
  40. steer,
  41. inject: () => {},
  42. cancel,
  43. runMaintenance: task => task(new AbortController().signal),
  44. whenIdle: () => Promise.resolve(),
  45. }
  46. ctx.agents.register(agent)
  47. ctx.provide('workspaceRegistry', { get: () => undefined, list: () => [] } as never)
  48. ctx.provide('agentDefaultModel', {
  49. currentSelection: () => ({ provider: 'fixture', model: 'fixture-model' }),
  50. saveSelection: () => Promise.resolve(),
  51. } as never)
  52. const selection: ModelSelectionRef = {
  53. current: { provider: 'fixture', model: 'fixture-model' },
  54. assembled: undefined,
  55. }
  56. const agents = {
  57. resolveAgent: () => Promise.resolve({ agent }),
  58. selectionFor: () => selection,
  59. serializeImageAdmission: <Value>(_agent: Agent, operation: () => Promise<Value>) => operation(),
  60. composeAgent: () => Promise.resolve({ setup: () => {} }),
  61. } as unknown as ApiSessionAgentController
  62. return { ctx, controller: new SessionCommandController(ctx, agents, '/workspace'), agent, inbox, steer, cancel }
  63. }
  64. async function expectFailure(operation: Promise<unknown>, code: string): Promise<void> {
  65. await expect(operation).rejects.toMatchObject({ code })
  66. }
  67. describe('Session queue commands', () => {
  68. it('edits, removes, steers, and rejects stale queue occurrences', async () => {
  69. const { ctx, controller, agent, inbox, steer, cancel } = await commandHarness()
  70. const queued = createUserMessage({ content: [{ type: 'text', text: 'queued' }], source: { kind: 'user' } })
  71. const nextStep = createUserMessage({ content: [{ type: 'text', text: 'step' }], source: { kind: 'user' } })
  72. inbox.append('next-turn', queued)
  73. inbox.append('next-step', nextStep)
  74. await expectFailure(Promise.resolve().then(() => controller.updateQueue({
  75. sessionId: agent.id,
  76. itemId: queued.id,
  77. action: {
  78. kind: 'edit',
  79. content: [{
  80. type: 'image',
  81. attachment: {
  82. attachmentId: AttachmentId('att-edit'), mediaType: 'image/png', bytes: 1, width: 1, height: 1,
  83. },
  84. }],
  85. },
  86. })), 'session/attachment-invalid')
  87. await expectFailure(Promise.resolve().then(() => controller.updateQueue({
  88. sessionId: SessionId('missing'), itemId: queued.id, action: { kind: 'remove' },
  89. })), 'session/queue-item-not-found')
  90. await expectFailure(Promise.resolve().then(() => controller.updateQueue({
  91. sessionId: agent.id, itemId: MessageId('missing'), action: { kind: 'remove' },
  92. })), 'session/queue-item-not-found')
  93. await expectFailure(Promise.resolve().then(() => controller.updateQueue({
  94. sessionId: agent.id, itemId: nextStep.id, action: { kind: 'steer' },
  95. })), 'session/steer-unavailable')
  96. Object.assign(agent, { status: 'idle' })
  97. await expectFailure(Promise.resolve().then(() => controller.updateQueue({
  98. sessionId: agent.id, itemId: queued.id, action: { kind: 'steer' },
  99. })), 'session/steer-unavailable')
  100. expect(controller.updateQueue({
  101. sessionId: agent.id,
  102. itemId: queued.id,
  103. action: { kind: 'edit', content: [{ type: 'text', text: 'edited' }] },
  104. })).toEqual({ accepted: true })
  105. expect(inbox.nextTurn[0]?.content).toEqual([{ type: 'text', text: 'edited' }])
  106. expect(controller.updateQueue({
  107. sessionId: agent.id, itemId: nextStep.id, action: { kind: 'remove' },
  108. })).toEqual({ accepted: true })
  109. Object.assign(agent, { status: 'running' })
  110. const steered = inbox.nextTurn[0]
  111. if (steered === undefined) throw new Error('missing edited queue item')
  112. expect(controller.updateQueue({
  113. sessionId: agent.id, itemId: steered.id, action: { kind: 'steer' },
  114. })).toEqual({ accepted: true })
  115. expect(steer).toHaveBeenCalledWith(steered)
  116. await expectFailure(Promise.resolve().then(() => controller.cancel({
  117. sessionId: SessionId('missing'),
  118. })), 'session/not-found')
  119. expect(controller.cancel({ sessionId: agent.id })).toEqual({ accepted: true })
  120. expect(cancel).toHaveBeenCalledWith({ kind: 'user' }, { keepInbox: true })
  121. await ctx.fiber.dispose()
  122. })
  123. })
  124. function imageRef(id: string): ImageAttachmentRef {
  125. return {
  126. attachmentId: AttachmentId(id),
  127. mediaType: 'image/png',
  128. bytes: 1,
  129. width: 1,
  130. height: 1,
  131. }
  132. }
  133. function event(type: string, seq: SessionSeq, data: unknown): SessionEvent {
  134. return { type, seq, time: seq + 1, data } as SessionEvent
  135. }
  136. async function persistedController(
  137. events: SessionEvent[],
  138. readImage: (ref: ImageAttachmentRef) => Promise<{ ref: ImageAttachmentRef; data: Uint8Array }>,
  139. ): Promise<{ ctx: Context; controller: SessionCommandController; sessionId: SessionId }> {
  140. const ctx = new Context()
  141. await ctx.plugin(SessionStore)
  142. const sessionId = SessionId('cold-attachment')
  143. const meta: SessionHeader = {
  144. version: 0,
  145. id: sessionId,
  146. createdAt: 1,
  147. cwd: '/workspace',
  148. isSeeded: false,
  149. }
  150. ctx.provide('sessionPersistence', testSessionPersistence(ctx, {
  151. list: () => Promise.resolve([meta]),
  152. inspect: () => Promise.resolve({
  153. meta,
  154. inheritedEventCount: SessionLogOffset(0),
  155. events,
  156. }),
  157. }) as never)
  158. installSessionReadTestServices(ctx)
  159. ctx.provide('attachments', { readImage } as never)
  160. const agents = { resolveAgent: vi.fn() } as unknown as ApiSessionAgentController
  161. return { ctx, controller: new SessionCommandController(ctx, agents, '/workspace'), sessionId }
  162. }
  163. describe('Session attachment authorization', () => {
  164. it('finds references in direct, message, inserted, nested, and streamed content', async () => {
  165. const nested = imageRef('nested')
  166. const message = imageRef('message')
  167. const inserted = imageRef('inserted')
  168. const streamed = imageRef('streamed')
  169. const events = [
  170. { ...event('fixture/direct', SessionSeq(0), {
  171. content: [null, [], { type: 'tool-result', content: [{ type: 'text', text: 'none' }] }, {
  172. type: 'tool-result', content: [{ type: 'image', attachment: nested }],
  173. }],
  174. }), ignorable: true as const },
  175. { ...event('assistant/message', SessionSeq(1), {
  176. turn: 1,
  177. step: 1,
  178. message: createAssistantMessage({
  179. content: [{ type: 'image', attachment: message }],
  180. source: { provider: 'fixture', model: 'fixture' },
  181. }),
  182. }), surfaceOp: 'append' as const },
  183. event('agent/inbox/spliced', SessionSeq(2), {
  184. target: 'next-turn',
  185. start: 0,
  186. inserted: [createUserMessage({
  187. content: [{ type: 'image', attachment: inserted }],
  188. source: { kind: 'user' },
  189. })],
  190. }),
  191. event('assistant/chunk', SessionSeq(3), {
  192. turn: 1,
  193. step: 1,
  194. chunk: { type: 'block-end', index: 0, block: { type: 'image', attachment: streamed } },
  195. }),
  196. ]
  197. const readImage = vi.fn((ref: ImageAttachmentRef) => Promise.resolve({ ref, data: Uint8Array.of(1) }))
  198. const { ctx, controller, sessionId } = await persistedController(events, readImage)
  199. for (const ref of [nested, message, inserted, streamed]) {
  200. await expect(controller.attachment({ sessionId, attachmentId: ref.attachmentId }))
  201. .resolves.toEqual({ attachment: ref, data: 'AQ==' })
  202. }
  203. expect(readImage).toHaveBeenCalledTimes(4)
  204. await ctx.fiber.dispose()
  205. })
  206. it('maps missing persistence identities and attachment backend failures', async () => {
  207. const noPersistence = new Context()
  208. await noPersistence.plugin(SessionStore)
  209. installSessionReadTestServices(noPersistence)
  210. const noPersistenceController = new SessionCommandController(
  211. noPersistence,
  212. { resolveAgent: vi.fn() } as unknown as ApiSessionAgentController,
  213. '/workspace',
  214. )
  215. await expectFailure(noPersistenceController.attachment({
  216. sessionId: SessionId('missing'), attachmentId: AttachmentId('att'),
  217. }), 'session/not-found')
  218. const missing = new Context()
  219. await missing.plugin(SessionStore)
  220. missing.provide('sessionPersistence', testSessionPersistence(missing, {
  221. list: () => Promise.resolve([]),
  222. inspect: vi.fn(),
  223. }) as never)
  224. installSessionReadTestServices(missing)
  225. const missingController = new SessionCommandController(
  226. missing,
  227. { resolveAgent: vi.fn() } as unknown as ApiSessionAgentController,
  228. '/workspace',
  229. )
  230. await expectFailure(missingController.attachment({
  231. sessionId: SessionId('missing'), attachmentId: 'att' as never,
  232. }), 'session/not-found')
  233. for (const thrown of [
  234. new AttachmentError('stored image is unavailable', 'ATTACHMENT_NOT_FOUND'),
  235. new Error('backend offline'),
  236. ]) {
  237. const ref = imageRef(`failure-${thrown.name}`)
  238. const fixture = await persistedController(
  239. [event('fixture/content', SessionSeq(0), { content: [{ type: 'image', attachment: ref }] })],
  240. () => Promise.reject(thrown),
  241. )
  242. await expectFailure(fixture.controller.attachment({
  243. sessionId: fixture.sessionId,
  244. attachmentId: ref.attachmentId,
  245. }), thrown instanceof AttachmentError ? 'session/attachment-invalid' : 'gateway/internal')
  246. await fixture.ctx.fiber.dispose()
  247. }
  248. })
  249. it('maps a cold observation failure to an internal authorization error', async () => {
  250. const ctx = new Context()
  251. await ctx.plugin(SessionStore)
  252. installSessionReadTestServices(ctx)
  253. vi.spyOn(ctx.sessionQuery, 'observeSession').mockRejectedValue(new Error('storage offline'))
  254. const controller = new SessionCommandController(
  255. ctx,
  256. { resolveAgent: vi.fn() } as unknown as ApiSessionAgentController,
  257. '/workspace',
  258. )
  259. await expectFailure(controller.attachment({
  260. sessionId: SessionId('unreadable'), attachmentId: AttachmentId('att'),
  261. }), 'gateway/internal')
  262. await ctx.fiber.dispose()
  263. })
  264. })