commands.ts 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617
  1. /** Session commands whose activation policy is explicit at each Remote method. */
  2. import { randomUUID } from 'node:crypto'
  3. import type { Context } from '@deepseek-ai/cordis'
  4. import { brandString } from '@deepseek-ai/dsh-brand'
  5. import type { Agent, ModelSelection as AgentModelSelection } from '@deepseek-ai/dsh-agent'
  6. import { AttachmentError } from '@deepseek-ai/dsh-attachment'
  7. import type {
  8. AttachmentAdmissionPart, FileAttachmentRef, ImageAttachmentRef,
  9. } from '@deepseek-ai/dsh-attachment'
  10. import type { FileUploadReceiptId } from '@deepseek-ai/dsh-client-file-upload/types'
  11. import type {} from '@deepseek-ai/dsh-client-file-upload'
  12. import {
  13. ReasoningEffortId, assistantStreamChunks, createUserMessage, freezeMessage,
  14. } from '@deepseek-ai/dsh-llm'
  15. import type { MessageSource } from '@deepseek-ai/dsh-llm'
  16. import { SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'
  17. import type { SessionEvent, SessionHeader, SessionId, UserMessage } from '@deepseek-ai/dsh-session'
  18. import { SessionQueryError, type SessionObservation } from '@deepseek-ai/dsh-session-query'
  19. import { SessionTitleInvalidError } from '@deepseek-ai/dsh-session-title'
  20. import { canonicalClientTimeZone } from '@deepseek-ai/dsh-util-time'
  21. import { RemoteError, remoteErrorOf } from '@deepseek-ai/dsh-typert-protocol'
  22. import type { Workspace } from '@deepseek-ai/dsh-workspace'
  23. import {
  24. ApiSessionAgentController,
  25. ApiSessionCwdConflict,
  26. ApiSessionNotFound,
  27. ApiSessionPresetConflict,
  28. ApiSessionSubagentOwnership,
  29. apiSessionSubagentOwnershipError,
  30. hasApiSessionSubagentOwner,
  31. inspectApiSession,
  32. } from './agent.ts'
  33. import type {
  34. SessionAttachmentRequest,
  35. SessionAttachmentValue,
  36. SessionCancelRequest,
  37. SessionCancelValue,
  38. SessionCreateRequest,
  39. SessionCreateValue,
  40. SessionForkRequest,
  41. SessionForkValue,
  42. SessionPromptRequest,
  43. SessionPromptValue,
  44. SessionRenameRequest,
  45. SessionRenameValue,
  46. SessionSelectModelRequest,
  47. SessionSelectModelValue,
  48. SessionUpdateQueueRequest,
  49. SessionUpdateQueueValue,
  50. SessionRequestId,
  51. } from './types.ts'
  52. interface SessionReadState {
  53. readonly id: SessionId
  54. readonly header: SessionHeader
  55. readonly events: readonly SessionEvent[]
  56. }
  57. /** Implements Session business commands delegated by the Session Controller Remote service. */
  58. export class SessionCommandController {
  59. /**
  60. * @param ctx - Host context carrying Agent, model, attachment, title, and Workspace services.
  61. * @param agents - sole owner of create, resume, and Session-local model selection.
  62. * @param defaultCwd - project directory used when create names neither a Workspace nor a cwd.
  63. */
  64. constructor(
  65. private readonly ctx: Context,
  66. private readonly agents: ApiSessionAgentController,
  67. private readonly defaultCwd: string,
  68. ) {}
  69. /**
  70. * Create or idempotently adopt one ordinary Session.
  71. * @param request - requested identity, location, and Agent preset.
  72. * @returns the Session identity and resolved preset when configured.
  73. */
  74. async create(request: SessionCreateRequest): Promise<SessionCreateValue> {
  75. if (request.workspaceId !== undefined && request.cwd !== undefined) {
  76. throw new RemoteError('gateway/bad-request', 'session.create accepts workspaceId or cwd, not both', {})
  77. }
  78. const sessionId = request.sessionId ?? brandString<SessionId>(`session-${randomUUID()}`)
  79. let workspace: Workspace | undefined
  80. if (request.workspaceId !== undefined) {
  81. workspace = this.ctx.workspaceRegistry.get(request.workspaceId)
  82. if (workspace === undefined) {
  83. throw new RemoteError('workspace/not-found', `workspace "${request.workspaceId}" not found`, {
  84. workspaceId: request.workspaceId,
  85. })
  86. }
  87. }
  88. const cwd = workspace?.path ?? request.cwd ?? this.defaultCwd
  89. let adopted: Agent
  90. try {
  91. adopted = await this.agents.ensureSession(
  92. sessionId,
  93. cwd,
  94. request.sessionId !== undefined,
  95. request.agentPreset,
  96. )
  97. } catch (error) {
  98. this.rejectCreation(sessionId, error)
  99. }
  100. if (workspace !== undefined) {
  101. try {
  102. await workspace.attachSession(sessionId)
  103. } catch (error) {
  104. throw new RemoteError(
  105. 'session/workspace-attach-failed',
  106. `session "${sessionId}" was created but could not attach to workspace "${workspace.id}": ${String(error)}`,
  107. { sessionId, workspaceId: workspace.id },
  108. )
  109. }
  110. }
  111. const agentPreset = this.agents.presetForSession(adopted.session)
  112. return { sessionId, ...(agentPreset === undefined ? {} : { agentPreset }) }
  113. }
  114. /**
  115. * Validate and install one Session-local model selection.
  116. * @param request - Session identity and requested model selection.
  117. * @returns the normalized selection installed for the Session.
  118. */
  119. async selectModel(request: SessionSelectModelRequest): Promise<SessionSelectModelValue> {
  120. const agent = await this.resolveAgent(request.sessionId)
  121. return this.agents.serializeImageAdmission(agent, async () => {
  122. try {
  123. const resolved = await this.ctx.llm.resolveCallConfig({
  124. provider: request.provider,
  125. model: request.model,
  126. ...(request.reasoningEffort === undefined
  127. ? {}
  128. : { reasoningEffort: ReasoningEffortId(request.reasoningEffort) }),
  129. })
  130. const selected: AgentModelSelection = {
  131. provider: resolved.provider,
  132. model: resolved.model,
  133. ...(resolved.reasoningEffort === undefined
  134. ? {}
  135. : { reasoningEffort: resolved.reasoningEffort }),
  136. }
  137. this.agents.selectForNextRequest(agent, selected)
  138. try {
  139. await this.ctx.agentDefaultModel.saveSelection(selected)
  140. } catch (error) {
  141. this.ctx.logger.warn(
  142. `session-controller: model selection changed for the Session but the default was not saved: ${String(error)}`,
  143. )
  144. }
  145. return { selected: { ...selected } }
  146. } catch (error) {
  147. if (remoteErrorOf(error) !== undefined) throw error
  148. throw new RemoteError(
  149. 'session/model-unavailable',
  150. error instanceof Error ? error.message : String(error),
  151. { provider: request.provider, model: request.model },
  152. )
  153. }
  154. })
  155. }
  156. /**
  157. * Normalize and append a user-owned Session title.
  158. * @param request - Session identity and proposed title.
  159. * @returns the accepted title and durable event sequence.
  160. */
  161. async rename(request: SessionRenameRequest): Promise<SessionRenameValue> {
  162. const agent = await this.resolveAgent(request.sessionId)
  163. const titles = this.ctx.get('sessionTitle')
  164. if (titles === undefined) {
  165. throw new RemoteError('gateway/internal', 'renaming is unavailable: this deployment mounts no session-title service', {})
  166. }
  167. try {
  168. const accepted = titles.rename(agent.session, request.title)
  169. return { title: accepted.title, seq: accepted.eventSeq }
  170. } catch (error) {
  171. if (error instanceof SessionTitleInvalidError) {
  172. throw new RemoteError('session/title-invalid', error.message, { sessionId: request.sessionId })
  173. }
  174. throw new RemoteError(
  175. 'gateway/internal',
  176. `failed to rename session "${request.sessionId}": ${String(error)}`,
  177. {},
  178. )
  179. }
  180. }
  181. /**
  182. * Create a new ordinary Session from one completed-turn prefix.
  183. * @param request - source Session and optional event anchor.
  184. * @returns the new Session identity.
  185. */
  186. async fork(request: SessionForkRequest): Promise<SessionForkValue> {
  187. let atSeq: ReturnType<typeof SessionSeq> | undefined
  188. try {
  189. atSeq = request.atSeq === undefined ? undefined : SessionSeq(request.atSeq)
  190. } catch {
  191. throw new RemoteError('gateway/bad-request', 'atSeq must be a non-negative safe integer', {})
  192. }
  193. let observed: SessionObservation
  194. try {
  195. observed = await this.ctx.sessionQuery.observeSession(request.sessionId)
  196. } catch (error) {
  197. if (error instanceof SessionQueryError
  198. && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') {
  199. throw new RemoteError('session/not-found', `session "${request.sessionId}" not found`, {
  200. sessionId: request.sessionId,
  201. })
  202. }
  203. throw new RemoteError(
  204. 'gateway/internal',
  205. `fork source unavailable for session "${request.sessionId}": ${String(error)}`,
  206. {},
  207. )
  208. }
  209. using source = observed
  210. const lastSeq = source.events.at(-1)?.seq ?? -1
  211. const anchoredBoundary = atSeq === undefined
  212. ? undefined
  213. : source.events.find(event => event.type === 'turn/end' && event.seq >= atSeq)
  214. const boundary = anchoredBoundary
  215. ?? (atSeq === undefined || atSeq > lastSeq
  216. ? source.events.findLast(event => event.type === 'turn/end')
  217. : undefined)
  218. if (boundary === undefined) {
  219. throw new RemoteError(
  220. 'session/fork-unavailable',
  221. atSeq !== undefined && atSeq <= lastSeq
  222. ? `session "${request.sessionId}" has not completed the turn containing event ${String(atSeq)}`
  223. : `session "${request.sessionId}" has no completed turn to fork from`,
  224. { sessionId: request.sessionId },
  225. )
  226. }
  227. let cut = SessionLogOffset(boundary.seq + 1)
  228. while (cut < source.events.length && source.events[cut]?.type !== 'turn/start') {
  229. cut = SessionLogOffset(cut + 1)
  230. }
  231. let workspace: Workspace | undefined
  232. try {
  233. workspace = await this.forkWorkspace(source.header)
  234. } catch (error) {
  235. throw new RemoteError(
  236. 'gateway/internal',
  237. `failed to resolve fork workspace for session "${request.sessionId}": ${String(error)}`,
  238. {},
  239. )
  240. }
  241. const childId = brandString<SessionId>(`session-${randomUUID()}`)
  242. const composition = await this.agents.composeAgent(this.agents.presetForObservation(source))
  243. try {
  244. const { provider, model } = this.ctx.agentDefaultModel.currentSelection()
  245. await this.ctx.agents.create({
  246. sessionId: childId,
  247. seed: source.events.slice(0, cut),
  248. inheritedEventCount: cut,
  249. meta: {
  250. ...(source.header.cwd === undefined ? {} : { cwd: source.header.cwd }),
  251. parentSession: source.header.id,
  252. isSeeded: true,
  253. ...(composition.agentPreset === undefined
  254. ? {}
  255. : { agentPreset: composition.agentPreset }),
  256. },
  257. agentOptions: { provider, model },
  258. setup: composition.setup,
  259. })
  260. } catch (error) {
  261. throw new RemoteError(
  262. 'gateway/internal',
  263. `failed to fork session "${request.sessionId}": ${String(error)}`,
  264. {},
  265. )
  266. }
  267. if (workspace !== undefined) {
  268. try {
  269. await workspace.attachSession(childId)
  270. } catch (error) {
  271. throw new RemoteError(
  272. 'session/workspace-attach-failed',
  273. `session "${childId}" was forked but could not attach to workspace "${workspace.id}": ${String(error)}`,
  274. { sessionId: childId, workspaceId: workspace.id },
  275. )
  276. }
  277. }
  278. return { sessionId: childId }
  279. }
  280. /**
  281. * Admit one browser prompt after explicit Agent resume and image validation.
  282. * @param request - Session identity, prompt content, source metadata, and delivery mode.
  283. * @returns acknowledgement that the Agent accepted the prompt.
  284. */
  285. async prompt(request: SessionPromptRequest): Promise<SessionPromptValue> {
  286. const clientTimeZone = request.clientTimeZone === undefined
  287. ? undefined
  288. : canonicalClientTimeZone(request.clientTimeZone)
  289. if (request.clientTimeZone !== undefined && clientTimeZone === undefined) {
  290. throw new RemoteError(
  291. 'session/invalid-time-zone',
  292. 'clientTimeZone must be UTC or a valid IANA Area/Location name',
  293. { value: request.clientTimeZone },
  294. )
  295. }
  296. const agent = await this.resolveAgent(request.sessionId)
  297. if (hasPromptRequest(agent, request.requestId)) return { accepted: true }
  298. const selection = this.agents.selectionFor(agent).current
  299. if (!routeServed(this.ctx, selection.provider)) {
  300. throw new RemoteError(
  301. 'session/model-unavailable',
  302. `no adapter serves provider "${selection.provider}"; select a model for this session`,
  303. { provider: selection.provider, model: selection.model },
  304. )
  305. }
  306. const source: MessageSource = {
  307. kind: 'user',
  308. rpcId: request.requestId,
  309. ...(clientTimeZone === undefined ? {} : { clientTimeZone }),
  310. }
  311. const hasImage = request.content.some(part => part.type === 'image')
  312. const admit = async (): Promise<SessionPromptValue> => {
  313. try {
  314. if (hasImage) {
  315. const current = this.agents.selectionFor(agent).current
  316. const model = await this.ctx.llm.resolveModelInfo(current.provider, current.model)
  317. if (model.inputModalities !== undefined && !model.inputModalities.includes('image')) {
  318. throw new RemoteError(
  319. 'session/attachment-invalid',
  320. `Model "${current.model}" does not support image input.`,
  321. { reason: 'MODEL_DOES_NOT_SUPPORT_IMAGES' },
  322. )
  323. }
  324. }
  325. const admission = resolvePromptFileReceipts(
  326. request.content,
  327. receiptId => this.ctx.fileUploads.resolve(agent, receiptId),
  328. )
  329. const content = await this.ctx.attachments.admitPromptContent(admission.content)
  330. const message: UserMessage = createUserMessage({ content, source })
  331. if (this.ctx.agents.get(agent.id) !== agent) {
  332. throw new RemoteError(
  333. 'session/not-found',
  334. `session "${agent.id}" was disposed during prompt admission`,
  335. { sessionId: agent.id },
  336. )
  337. }
  338. using binding = this.ctx.fileUploads.bindPrompt(agent, admission.receiptIds, request.requestId)
  339. if (request.mode === 'steer') agent.steer(message)
  340. else agent.followup(message)
  341. binding.commit()
  342. } catch (error) {
  343. if (remoteErrorOf(error) !== undefined) throw error
  344. if (error instanceof AttachmentError) {
  345. throw new RemoteError('session/attachment-invalid', error.message, { reason: error.code })
  346. }
  347. throw new RemoteError('session/agent-busy', 'prompt rejected', { reason: String(error) })
  348. }
  349. return { accepted: true }
  350. }
  351. return hasImage ? this.agents.serializeImageAdmission(agent, admit) : admit()
  352. }
  353. /**
  354. * Read one durable image after proving the Session log references it.
  355. * @param request - Session and attachment identities used for authorization.
  356. * @returns the durable attachment reference and base64-encoded bytes.
  357. */
  358. async attachment(request: SessionAttachmentRequest): Promise<SessionAttachmentValue> {
  359. let source: SessionReadState
  360. try {
  361. source = await this.readSessionState(request.sessionId)
  362. } catch (error) {
  363. if (error instanceof ApiSessionNotFound) {
  364. throw new RemoteError('session/not-found', error.message, { sessionId: request.sessionId })
  365. }
  366. throw new RemoteError(
  367. 'gateway/internal',
  368. `attachment authorization unavailable for session "${request.sessionId}": ${String(error)}`,
  369. {},
  370. )
  371. }
  372. const ref = referencedImage(source.events, String(request.attachmentId))
  373. if (ref === undefined) {
  374. throw new RemoteError(
  375. 'session/attachment-invalid',
  376. 'Image is not referenced by this session.',
  377. { reason: 'ATTACHMENT_NOT_REFERENCED' },
  378. )
  379. }
  380. try {
  381. const stored = await this.ctx.attachments.readImage(ref)
  382. return {
  383. attachment: stored.ref,
  384. data: Buffer.from(stored.data).toString('base64'),
  385. }
  386. } catch (error) {
  387. if (error instanceof AttachmentError) {
  388. throw new RemoteError('session/attachment-invalid', error.message, { reason: error.code })
  389. }
  390. throw new RemoteError('gateway/internal', 'Unable to read image attachment.', {})
  391. }
  392. }
  393. /**
  394. * Mutate one still-pending queue occurrence without resuming a cold Agent.
  395. * @param request - Session, queue item, and requested mutation.
  396. * @returns acknowledgement that the queue mutation was applied.
  397. */
  398. updateQueue(request: SessionUpdateQueueRequest): SessionUpdateQueueValue {
  399. if (request.action.kind === 'edit'
  400. && request.action.content.some(block => block.type !== 'text')) {
  401. throw new RemoteError(
  402. 'session/attachment-invalid',
  403. 'queue edits accept text content only',
  404. { reason: 'QUEUE_EDIT_NON_TEXT' },
  405. )
  406. }
  407. const agent = this.ctx.agents.get(request.sessionId)
  408. if (agent !== undefined && hasApiSessionSubagentOwner(this.ctx, agent.session, agent)) {
  409. throw apiSessionSubagentOwnershipError(request.sessionId)
  410. }
  411. if (agent === undefined) {
  412. throw new RemoteError('session/queue-item-not-found', 'queued item is no longer pending', { itemId: request.itemId })
  413. }
  414. const nextTurn = agent.inbox.nextTurn.find(message => message.id === request.itemId)
  415. const nextStep = agent.inbox.nextStep.find(message => message.id === request.itemId)
  416. const located = nextTurn === undefined
  417. ? nextStep === undefined ? undefined : { target: 'next-step' as const, message: nextStep }
  418. : { target: 'next-turn' as const, message: nextTurn }
  419. if (located === undefined) {
  420. throw new RemoteError('session/queue-item-not-found', 'queued item is no longer pending', { itemId: request.itemId })
  421. }
  422. const { target, message } = located
  423. if (request.action.kind === 'steer' && (target !== 'next-turn' || agent.status !== 'running')) {
  424. throw new RemoteError('session/steer-unavailable', 'current turn no longer accepts steering', { itemId: request.itemId })
  425. }
  426. if (request.action.kind === 'edit') {
  427. agent.inbox.replace(request.itemId, freezeMessage<UserMessage>({
  428. ...message,
  429. content: [...request.action.content],
  430. }))
  431. } else {
  432. agent.inbox.remove(request.itemId)
  433. if (request.action.kind === 'remove') {
  434. const source = message.source
  435. if (source.kind === 'user' && 'rpcId' in source) {
  436. this.ctx.fileUploads.retirePrompt(agent, source.rpcId)
  437. }
  438. }
  439. if (request.action.kind === 'steer') agent.steer(message)
  440. }
  441. return { accepted: true }
  442. }
  443. /**
  444. * Cancel one live ordinary Agent while retaining pending inbox work.
  445. * @param request - Session whose active Agent turn is cancelled.
  446. * @returns acknowledgement that cancellation was requested.
  447. */
  448. cancel(request: SessionCancelRequest): SessionCancelValue {
  449. const agent = this.ctx.agents.get(request.sessionId)
  450. if (agent === undefined) {
  451. throw new RemoteError(
  452. 'session/not-found',
  453. `session "${request.sessionId}" not found (not attached)`,
  454. { sessionId: request.sessionId },
  455. )
  456. }
  457. if (hasApiSessionSubagentOwner(this.ctx, agent.session, agent)) {
  458. throw apiSessionSubagentOwnershipError(request.sessionId)
  459. }
  460. agent.cancel({ kind: 'user' }, { keepInbox: true })
  461. return { accepted: true }
  462. }
  463. private async resolveAgent(sessionId: SessionId): Promise<Agent> {
  464. const found = await this.agents.resolveAgent(sessionId)
  465. if ('error' in found) throw found.error
  466. return found.agent
  467. }
  468. private rejectCreation(sessionId: SessionId, error: unknown): never {
  469. if (remoteErrorOf(error) !== undefined) throw error
  470. if (error instanceof ApiSessionPresetConflict) {
  471. throw new RemoteError('agent-preset/conflict', error.message, {
  472. sessionId: error.sessionId,
  473. requestedPreset: error.requestedPreset,
  474. ...(error.existingPreset === undefined ? {} : { existingPreset: error.existingPreset }),
  475. })
  476. }
  477. if (error instanceof ApiSessionCwdConflict) {
  478. throw new RemoteError('session/conflict', error.message, {
  479. sessionId: error.sessionId,
  480. requestedCwd: error.requestedCwd,
  481. ...(error.existingCwd === undefined ? {} : { existingCwd: error.existingCwd }),
  482. })
  483. }
  484. if (error instanceof ApiSessionSubagentOwnership) {
  485. throw apiSessionSubagentOwnershipError(error.sessionId)
  486. }
  487. throw new RemoteError('gateway/internal', `failed to create session "${sessionId}": ${String(error)}`, {})
  488. }
  489. private async readSessionState(sessionId: SessionId): Promise<SessionReadState> {
  490. const attached = this.ctx.sessions.get(sessionId)
  491. if (attached !== undefined) {
  492. return { id: attached.id, header: attached.header, events: attached.snapshotEvents() }
  493. }
  494. const inspected = await inspectApiSession(this.ctx, sessionId)
  495. return { id: inspected.meta.id, header: inspected.meta, events: inspected.events }
  496. }
  497. private async forkWorkspace(source: SessionHeader): Promise<Workspace | undefined> {
  498. const workspaces = this.ctx.workspaceRegistry.list()
  499. const direct = workspaces.find(workspace => workspace.sessionIds.includes(source.id))
  500. if (direct !== undefined || source.origin !== 'subagent') return direct
  501. const lineage = await this.ctx.sessionQuery.traceSession(source.id)
  502. for (const ancestor of lineage.ancestors) {
  503. const workspace = workspaces.find(candidate => candidate.sessionIds.includes(ancestor.header.id))
  504. if (workspace !== undefined) return workspace
  505. }
  506. return undefined
  507. }
  508. }
  509. function resolvePromptFileReceipts(
  510. content: SessionPromptRequest['content'],
  511. stagedFile: (receiptId: FileUploadReceiptId) => FileAttachmentRef | undefined,
  512. ): { readonly content: AttachmentAdmissionPart[]; readonly receiptIds: readonly FileUploadReceiptId[] } {
  513. const receiptIds = new Set<FileUploadReceiptId>()
  514. const resolved = content.map((part): AttachmentAdmissionPart => {
  515. if (part.type !== 'file') return part
  516. const attachment = stagedFile(part.receiptId)
  517. if (attachment === undefined) {
  518. throw new RemoteError(
  519. 'session/attachment-invalid',
  520. 'File was not uploaded for this session.',
  521. { reason: 'FILE_NOT_STAGED' },
  522. )
  523. }
  524. receiptIds.add(part.receiptId)
  525. return { type: 'file', attachment }
  526. })
  527. return { content: resolved, receiptIds: [...receiptIds] }
  528. }
  529. function hasPromptRequest(agent: Agent, requestId: SessionRequestId): boolean {
  530. const matches = (message: UserMessage): boolean => {
  531. const source = message.source
  532. return source.kind === 'user' && 'rpcId' in source && source.rpcId === requestId
  533. }
  534. if (agent.inbox.nextTurn.some(matches) || agent.inbox.nextStep.some(matches)) return true
  535. return agent.session.snapshotEvents().some((event) => {
  536. if (event.type !== 'user/message') return false
  537. const source = event.data.source
  538. return source.kind === 'user' && 'rpcId' in source && source.rpcId === requestId
  539. })
  540. }
  541. function imageBlockIn(
  542. content: unknown,
  543. match: (ref: ImageAttachmentRef) => boolean,
  544. ): ImageAttachmentRef | undefined {
  545. if (!Array.isArray(content)) return undefined
  546. for (const value of content) {
  547. if (typeof value !== 'object' || value === null || Array.isArray(value)) continue
  548. const block = value as { readonly type?: unknown; readonly attachment?: unknown; readonly content?: unknown }
  549. if (block.type === 'image' && typeof block.attachment === 'object' && block.attachment !== null) {
  550. const ref = block.attachment as ImageAttachmentRef
  551. if (match(ref)) return ref
  552. }
  553. if (block.type === 'tool-result') {
  554. const nested = imageBlockIn(block.content, match)
  555. if (nested !== undefined) return nested
  556. }
  557. }
  558. return undefined
  559. }
  560. function imageInEvent(
  561. event: SessionEvent,
  562. match: (ref: ImageAttachmentRef) => boolean,
  563. ): ImageAttachmentRef | undefined {
  564. const data = event.data as {
  565. readonly content?: unknown
  566. readonly message?: { readonly content?: unknown }
  567. readonly inserted?: readonly { readonly content?: unknown }[]
  568. }
  569. const direct = imageBlockIn(data.content, match)
  570. if (direct !== undefined) return direct
  571. const message = imageBlockIn(data.message?.content, match)
  572. if (message !== undefined) return message
  573. for (const inserted of data.inserted ?? []) {
  574. const found = imageBlockIn(inserted.content, match)
  575. if (found !== undefined) return found
  576. }
  577. if (event.type === 'assistant/message' || event.type === 'assistant/attempt') {
  578. for (const chunk of assistantStreamChunks(event.data.stream, 'block-end')) {
  579. const found = imageBlockIn([chunk.block], match)
  580. if (found !== undefined) return found
  581. }
  582. }
  583. return undefined
  584. }
  585. function referencedImage(
  586. events: readonly SessionEvent[],
  587. attachmentId: string,
  588. ): ImageAttachmentRef | undefined {
  589. for (const event of events) {
  590. const found = imageInEvent(event, ref => String(ref.attachmentId) === attachmentId)
  591. if (found !== undefined) return found
  592. }
  593. return undefined
  594. }
  595. function routeServed(ctx: Context, provider: string): boolean {
  596. return ctx.llm.listProviders().some(entry => entry.id === provider)
  597. }