commands.ts 26 KB

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