api-proxy.ts 71 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630
  1. /**
  2. * Host-side ApiProxy implementation. Signature discipline: unary takes the
  3. * narrow RpcRequest<P> and echoes request.rpcId on the RpcResponse<T>.
  4. */
  5. import { randomUUID } from 'node:crypto'
  6. import { mkdir, stat } from 'node:fs/promises'
  7. import { join } from 'node:path'
  8. import type { Context } from 'cordis'
  9. import { installAgentLlmTarget } from '@deepseek-ai/dsh-agent'
  10. import type {
  11. Agent, AgentLlmTarget, AgentLlmTargetRef, AgentStatus, InboxPlacement,
  12. } from '@deepseek-ai/dsh-agent'
  13. import { createUserMessage, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  14. import { errorChain } from '@deepseek-ai/dsh-llm'
  15. import type { MessageId, MessageSource } from '@deepseek-ai/dsh-llm'
  16. import { isAppendSurfaceEvent } from '@deepseek-ai/dsh-session'
  17. import type { Session, SessionEvent, SessionHeader, SessionId, UserMessage } from '@deepseek-ai/dsh-session'
  18. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  19. import type { Workspace, WorkspaceRecord } from '@deepseek-ai/dsh-workspace'
  20. import {
  21. workspaceDomainState, workspaceRecord, WorkspaceId as brandWorkspaceId,
  22. WorkspaceMoveInvalidError, WorkspaceNameConflictError,
  23. } from '@deepseek-ai/dsh-workspace'
  24. // Type-only: brings the `ctx.tools` Context merge into this program (viewFor reads presenters).
  25. import type {} from '@deepseek-ai/dsh-tools'
  26. import type {
  27. ApiProxy, GoalRef, HistoryEntry, HostFrame, ModelCatalogFailure, ModelProviderGroup, ModelReasoning,
  28. MuxFrame, QuestionResponsePayload, SessionProjectionsBlock, SessionSummary, ToolEventView,
  29. WorkspaceId, WorkspaceView,
  30. } from './api/index.ts'
  31. // Type-only: resolves `ctx.get('sessionProjections')` to the projection registry.
  32. import type {} from '@deepseek-ai/dsh-session-projection'
  33. // Type-only: resolves `ctx.get('sessionProjectionCache')` (the cold listing column).
  34. import type {} from '@deepseek-ai/dsh-session-projection-cache'
  35. // GoalError narrows domain rejections to their stable codes at the wire boundary.
  36. import { GoalError } from '@deepseek-ai/dsh-goal'
  37. import type { GoalRef as CoreGoalRef } from '@deepseek-ai/dsh-goal'
  38. // Type-only edges: resolve `ctx.get('commands')`, the `commands/change` event, and `ctx.get('skills')`.
  39. import type {} from '@deepseek-ai/dsh-commands'
  40. import type {} from '@deepseek-ai/dsh-skill'
  41. import type { CallId } from '@deepseek-ai/dsh-llm/brand'
  42. import type { ApprovalOutcome, ApprovalRequestId } from '@deepseek-ai/dsh-user-approval'
  43. // Side-effect type import: resolves the `approval/request` waterfall and
  44. // `ctx.get('approval')` without a value dependency on the seam (optional composition).
  45. import type {} from '@deepseek-ai/dsh-user-approval'
  46. import { approvalResponsePayloadSchema } from './api/approvals.schema.ts'
  47. import { questionResponsePayloadSchema } from './api/questions.schema.ts'
  48. import type { ClientResponse, RpcError, RpcReceipt, RpcRequest, RpcResponse } from './api/rpc.ts'
  49. import { RpcId } from './api/rpc.ts'
  50. import type {
  51. AskUserQuestionAnswer, AskUserQuestionItem, AskUserQuestionRequest,
  52. } from '@deepseek-ai/dsh-user-interaction'
  53. import { UserInteractionError } from '@deepseek-ai/dsh-user-interaction'
  54. import { DirectoryPickerError } from '@deepseek-ai/dsh-host-directory-picker'
  55. import { openNativePath } from './native-path-opener.ts'
  56. /** Page size when history is called without maxMessages. */
  57. const DEFAULT_MAX_MESSAGES = 50
  58. /** Conversation message event types (the pagination counting unit). */
  59. const MESSAGE_TYPES = new Set(['user/message', 'assistant/message', 'steering/message'])
  60. /**
  61. * Message-boundary pagination: count maxMessages append-origin messages
  62. * backwards from the window tail. Replacement copies never entered the
  63. * conversation a reader sees — they restate a shadowed range for the model
  64. * alone — so they consume no quota; the page stays one contiguous raw range,
  65. * which keeps a compaction's log-only provenance on the same page as its
  66. * replacement. The cut is the starting seq of the oldest message group (chunks
  67. * group via sourceEventSeqs — never cut mid-message). The tail page naturally
  68. * includes the in-progress partial.
  69. */
  70. function paginate(
  71. events: readonly SessionEvent[],
  72. beforeSeq: number | undefined,
  73. maxMessages: number,
  74. ): { events: SessionEvent[]; hasMore: boolean } {
  75. const window = beforeSeq === undefined ? [...events] : events.filter(event => event.seq < beforeSeq)
  76. let count = 0
  77. let cut = 0
  78. for (let i = window.length - 1; i >= 0; i--) {
  79. const event = window[i] as SessionEvent
  80. if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue
  81. count++
  82. const sources = (event as { sourceEventSeqs?: number[] }).sourceEventSeqs
  83. const groupStart = sources !== undefined && sources.length > 0 ? Math.min(event.seq, ...sources) : event.seq
  84. if (count >= maxMessages) {
  85. cut = groupStart
  86. break
  87. }
  88. }
  89. const page = window.filter(event => event.seq >= cut)
  90. return { events: page, hasMore: cut > 0 }
  91. }
  92. /** Wrap an ok result echoing the request's rpcId. */
  93. function ok<T>(request: RpcRequest<unknown>, value: T): RpcResponse<T> {
  94. return { rpcId: request.rpcId, result: { ok: true, value } }
  95. }
  96. /** Wrap an error result echoing the request's rpcId. */
  97. function err<T>(request: RpcRequest<unknown>, error: RpcError): RpcResponse<T> {
  98. return { rpcId: request.rpcId, result: { ok: false, error } }
  99. }
  100. /** Simple async queue: core callbacks push, the AsyncIterable pulls; abort/return cleans up. */
  101. class FrameQueue<F> {
  102. private buffer: F[] = []
  103. private waiter: (() => void) | undefined
  104. private done = false
  105. push(item: F): void {
  106. if (this.done) return
  107. this.buffer.push(item)
  108. this.waiter?.()
  109. }
  110. end(): void {
  111. this.done = true
  112. this.waiter?.()
  113. }
  114. async *iterate(signal: AbortSignal, cleanup: () => void): AsyncGenerator<F> {
  115. const onAbort = (): void => { this.end() }
  116. signal.addEventListener('abort', onAbort, { once: true })
  117. try {
  118. while (true) {
  119. while (this.buffer.length > 0) yield this.buffer.shift() as F
  120. if (this.done || signal.aborted) return
  121. await new Promise<void>((resolve) => { this.waiter = resolve })
  122. this.waiter = undefined
  123. }
  124. } finally {
  125. signal.removeEventListener('abort', onAbort)
  126. cleanup()
  127. }
  128. }
  129. }
  130. /**
  131. * Server-side frame mint: pure pushes get a fresh rpcId per frame (answerable
  132. * frames — approval/question requested — mint their stable id in their
  133. * pending registries instead).
  134. */
  135. function frame<F>(payload: F): RpcRequest<F> {
  136. return { rpcId: RpcId(randomUUID()), payload }
  137. }
  138. /** Queue the subscription baseline frame. */
  139. function subscribeSession(queue: FrameQueue<RpcRequest<MuxFrame>>, session: Session): void {
  140. queue.push(frame({ type: 'session/subscribed', sessionId: session.id, lastSeq: session.seq - 1 }))
  141. }
  142. /**
  143. * Whether the session's conversation has started: no turn has run yet (a
  144. * turn is one model-loop execution). Standalone plugin events — command
  145. * lifecycle records, plan/mode, titles, goals — never open a turn, so
  146. * running `/plan` or `/goal` on a fresh session keeps it blank
  147. * (list-hidden, reusable).
  148. */
  149. function sessionBlank(session: Session): boolean {
  150. return !session.events.some(event => event.type === 'turn/start')
  151. }
  152. /** SessionSummary projection for attached (in-memory) sessions. */
  153. function summarize(session: Session, running: boolean): SessionSummary {
  154. return {
  155. sessionId: session.id,
  156. updatedAt: session.events.at(-1)?.time ?? session.header.createdAt,
  157. running,
  158. blank: sessionBlank(session),
  159. ...session.header.parentSession === undefined ? {} : { parentSessionId: session.header.parentSession },
  160. ...session.header.cwd === undefined ? {} : { cwd: session.header.cwd },
  161. }
  162. }
  163. /**
  164. * SessionSummary projection for cold (persisted, unattached) sessions.
  165. * updatedAt is the log file's mtime; backends without a per-session file
  166. * (locate() undefined) fall back to the header's createdAt.
  167. */
  168. async function summarizeCold(persistence: SessionPersistence, meta: SessionHeader): Promise<SessionSummary> {
  169. let updatedAt = meta.createdAt
  170. const location = persistence.locate(meta)
  171. if (location !== undefined) {
  172. try {
  173. updatedAt = (await stat(location.path)).mtimeMs
  174. } catch {
  175. // The log vanished between list() and stat() (concurrent cleanup); createdAt stands in.
  176. }
  177. }
  178. return {
  179. sessionId: meta.id,
  180. updatedAt,
  181. running: false,
  182. // Lazy persistence keeps never-appended sessions out of list(); reading
  183. // a cold log to check for turns would defeat the index read, so a listed
  184. // cold session is served as not-blank (its log holds its conversation).
  185. blank: false,
  186. ...meta.parentSession === undefined ? {} : { parentSessionId: meta.parentSession },
  187. /* v8 ignore next -- the empty arm needs a cwd-less meta, but list()
  188. filters those out (legacy logs are not served); the conditional mirrors
  189. summarize() shape. */
  190. ...meta.cwd === undefined ? {} : { cwd: meta.cwd },
  191. }
  192. }
  193. /** Map a browse-primitive failure onto the wire error vocabulary (unknown throws stay internal). */
  194. function directoryError(error: unknown): RpcError {
  195. if (error instanceof DirectoryPickerError) {
  196. return { code: error.code, message: error.message, details: { path: error.path } }
  197. }
  198. return { code: 'internal', message: error instanceof Error ? error.message : String(error), details: {} }
  199. }
  200. /** Resolved Host routing and project-directory defaults consumed by the API implementation. */
  201. export interface ApiProxyDefaults {
  202. provider: string
  203. model: string
  204. /** Default project directory for new sessions whose create request carries no cwd. */
  205. cwd: string
  206. /** Parent directory for name-created workspaces. */
  207. workspaceRoot: string
  208. /** Native open-with-default-application; injectable for carrier tests. */
  209. openPath?: (path: string, signal: AbortSignal) => Promise<void>
  210. }
  211. /** The tool/call payload fields the presenter path reads. */
  212. interface ToolCallData { callId: string; name: string; arguments: string }
  213. /**
  214. * One outstanding approval question: the stable server-request id, the frame
  215. * material replayed to late mux subscribers, and the resolver that settles the
  216. * answerer's promise back into `ctx.approval`.
  217. */
  218. interface PendingApproval {
  219. rpcId: RpcId
  220. sessionId: SessionId
  221. approvalId: ApprovalRequestId
  222. toolName: string
  223. callId?: CallId
  224. reason?: string
  225. resolve(outcome: ApprovalOutcome): void
  226. }
  227. /** Project a pending entry into its answerable mux frame (initial push and mux-open replay share it). */
  228. function requestedFrame(pending: PendingApproval): RpcRequest<MuxFrame> {
  229. return {
  230. rpcId: pending.rpcId,
  231. payload: {
  232. type: 'approval/requested',
  233. sessionId: pending.sessionId,
  234. approvalId: pending.approvalId,
  235. toolName: pending.toolName,
  236. ...pending.callId === undefined ? {} : { callId: pending.callId },
  237. ...pending.reason === undefined ? {} : { reason: pending.reason },
  238. },
  239. }
  240. }
  241. /** One host-owned question wait, addressed by the stable server-request id. */
  242. interface PendingQuestion {
  243. rpcId: RpcId
  244. sessionId: SessionId
  245. questions: AskUserQuestionItem[]
  246. resolve: (answer: AskUserQuestionAnswer) => void
  247. reject: (error: UserInteractionError) => void
  248. signal?: AbortSignal
  249. onAbort?: () => void
  250. }
  251. /** Validate one answer batch against the exact question request it resolves. */
  252. function matchesQuestions(payload: QuestionResponsePayload, pending: PendingQuestion): boolean {
  253. if (payload.sessionId !== pending.sessionId) return false
  254. const answers = payload.answer.answers
  255. if (answers.length !== pending.questions.length) return false
  256. return answers.every((answer, index) => {
  257. const question = pending.questions[index] as AskUserQuestionItem
  258. if (answer.id !== question.id) return false
  259. if (new Set(answer.selected).size !== answer.selected.length) return false
  260. const custom = answer.custom?.trim()
  261. if (custom !== undefined && custom === '') return false
  262. if (custom !== undefined && answer.selected.length > 0) return false
  263. if (question.multiSelect !== true && answer.selected.length > 1) return false
  264. const labels = new Set(question.options?.map(option => option.label) ?? [])
  265. return answer.selected.every(label => labels.has(label))
  266. })
  267. }
  268. /**
  269. * Compute the render intent for a tool/call or tool/result event through the
  270. * presenters registered at this moment; every other event type gets none. A
  271. * result's presenter needs its call's parsed args — `argsFor` supplies them
  272. * (live: the per-session call table; history: an in-page backscan), returning
  273. * undefined when the pairing is unavailable (e.g. the call fell off the page),
  274. * which soft-falls to no view. Presenter or JSON.parse throws also soft-fall:
  275. * the client's documented default (generic JSON card) covers every miss.
  276. */
  277. function viewFor(ctx: Context, event: SessionEvent, argsFor: (callId: string) => unknown): ToolEventView | undefined {
  278. try {
  279. if (event.type === 'tool/call') {
  280. const { name, arguments: raw } = event.data as ToolCallData
  281. const view = ctx.tools.get(name)?.presentCall?.(JSON.parse(raw))
  282. return view === undefined ? undefined : { for: 'call', view }
  283. }
  284. if (event.type === 'tool/result') {
  285. const { message, meta } = event.data
  286. const [result] = message.content
  287. const callId = message.source.callId
  288. const call = argsFor(callId) as { name: string; args: unknown } | undefined
  289. if (call === undefined) return undefined
  290. const view = ctx.tools.get(call.name)?.presentResult?.(call.args, {
  291. content: result.content,
  292. isError: result.isError === true,
  293. ...meta === undefined ? {} : { meta },
  294. })
  295. return view === undefined ? undefined : { for: 'result', view }
  296. }
  297. } catch (error: unknown) {
  298. // A throwing presenter (or unparseable arguments) must not break delivery;
  299. // the event still ships, just without a view.
  300. console.error(`api-proxy: presenter failed for ${event.type}, falling back to generic: ${String(error)}`)
  301. }
  302. return undefined
  303. }
  304. /**
  305. * Resolve a tool/result's call pairing by scanning a window of events backwards
  306. * for the matching tool/call. Used by the history path (the page is the
  307. * window — a cross-page pairing soft-falls to no view) and by live-path table
  308. * misses after a reconnect-eviction.
  309. */
  310. function backscanArgs(events: readonly SessionEvent[], callId: string): { name: string; args: unknown } | undefined {
  311. for (let i = events.length - 1; i >= 0; i--) {
  312. const event = events[i] as SessionEvent
  313. if (event.type !== 'tool/call') continue
  314. const data = event.data as ToolCallData
  315. if (data.callId !== callId) continue
  316. try {
  317. return { name: data.name, args: JSON.parse(data.arguments) }
  318. } catch {
  319. // Unparseable stored arguments: same soft-fall as a live parse failure.
  320. return undefined
  321. }
  322. }
  323. return undefined
  324. }
  325. /**
  326. * The projection baseline for one history tail page: the registry's
  327. * watermark-cache snapshot — one fully synchronous read (no await between the
  328. * page slice and this), so all values and `asOfSeq` form a single consistent
  329. * cut and `asOfSeq` equals the window tail event seq. The carrier holds zero
  330. * domain knowledge (each value passed its unit's own schema inside the
  331. * registry). An absent registry means the deployment has no projection seam:
  332. * the whole block is absent and clients treat every key as capability-absent.
  333. */
  334. function projectionsFor(ctx: Context, agent: Agent): SessionProjectionsBlock | undefined {
  335. const registry = ctx.get('sessionProjections')
  336. if (registry === undefined) return undefined
  337. return registry.snapshot(agent.session)
  338. }
  339. /**
  340. * The projection baseline of one session.list row, fail-soft: attached
  341. * sessions cut the registry's live watermark cache; cold sessions view the
  342. * persisted projection cache's identity-checked stored rows (zero log loads
  343. * either way — the listing use case the cache exists for). The block shape
  344. * (values + asOfSeq) matches the history tail's, so a client seeds its
  345. * value store under the same higher-seq-wins rule. Any failure — and an
  346. * empty value set — yields an absent block: a listing without projections
  347. * is degraded, never broken.
  348. */
  349. function listProjectionsFor(ctx: Context, meta: SessionHeader, session: Session | undefined): SessionProjectionsBlock | undefined {
  350. try {
  351. const block = session !== undefined
  352. ? ctx.get('sessionProjections')?.snapshot(session)
  353. : ctx.get('sessionProjectionCache')?.cachedSnapshot(meta)
  354. return block !== undefined && Object.keys(block.values).length > 0 ? block : undefined
  355. } catch (error) {
  356. ctx.logger.warn(`session.list: projection column for "${meta.id}" failed (serving the row without it): ${String(error)}`)
  357. return undefined
  358. }
  359. }
  360. /**
  361. * Thrown by the cold-resume path when the id names no servable session
  362. * (absent from the store, or a pre-project legacy log without a cwd).
  363. */
  364. class SessionNotFound extends Error {}
  365. /** Requested identity already belongs to a session with another project cwd. */
  366. class SessionCwdConflict extends Error {
  367. constructor(
  368. readonly sessionId: SessionId,
  369. readonly requestedCwd: string,
  370. readonly existingCwd: string | undefined,
  371. ) {
  372. super(
  373. `session "${sessionId}" already exists with cwd ${JSON.stringify(existingCwd)}; `
  374. + `requested ${JSON.stringify(requestedCwd)}`,
  375. )
  376. }
  377. }
  378. /** Host failed before the registry could adopt a name-created directory. */
  379. class WorkspaceDirectoryCreationError extends Error {}
  380. /** Shared workspace-not-found error response of the workspace.* mutation rows. */
  381. function workspaceNotFound<T>(request: RpcRequest<unknown>, workspaceId: string): RpcResponse<T> {
  382. return err(request, {
  383. code: 'workspace-not-found',
  384. message: `workspace "${workspaceId}" not found`,
  385. details: { workspaceId },
  386. })
  387. }
  388. /** Wire projection of one workspace entity (the workspace.* value row). */
  389. function workspaceView(workspace: Workspace): WorkspaceView {
  390. return {
  391. workspaceId: workspace.id,
  392. path: workspace.path,
  393. title: workspace.title,
  394. sessionIds: [...workspace.sessionIds],
  395. createdAt: workspace.createdAt,
  396. updatedAt: workspace.updatedAt,
  397. }
  398. }
  399. /** Wire projection of the durable record carried by `domain/changed`. */
  400. function changedWorkspaceView(workspaceId: string, value: unknown): WorkspaceView {
  401. const record: WorkspaceRecord = workspaceRecord.parse(value)
  402. return {
  403. workspaceId: workspaceId as WorkspaceId,
  404. path: record.path,
  405. title: record.title,
  406. sessionIds: [...record.sessionIds],
  407. createdAt: record.createdAt,
  408. updatedAt: record.updatedAt,
  409. }
  410. }
  411. /**
  412. * Implement ApiProxy over a composed host context.
  413. * @param ctx - a context with the Host spine and Workspace registry mounted.
  414. * @param defaults - host routing and project-directory defaults.
  415. * @returns the ApiProxy implementation.
  416. */
  417. export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiProxy {
  418. const agentOptions = { provider: defaults.provider, model: defaults.model }
  419. type WebLlmTargetRef = AgentLlmTargetRef & { current: AgentLlmTarget }
  420. const targets = new WeakMap<Agent, WebLlmTargetRef>()
  421. /** Implicit resume of cold sessions, deduplicating concurrent calls (follows the jsonrpc sessionCreations precedent). */
  422. const resumes = new Map<SessionId, Promise<Agent>>()
  423. /** Client-chosen identity creation/resume, deduplicated across concurrent retries. */
  424. const sessionCreations = new Map<SessionId, Promise<Agent>>()
  425. /** Serializes path ownership checks with record creation across spellings. */
  426. let workspaceCreationChain = Promise.resolve()
  427. const pendingQuestions = new Map<RpcId, PendingQuestion>()
  428. const pendingApprovals = new Map<RpcId, PendingApproval>()
  429. const muxQueues = new Set<FrameQueue<RpcRequest<MuxFrame>>>()
  430. /**
  431. * Install or return the session-local target that prompt assembly snapshots.
  432. * Seed order: latest logged request/header, else the host default routing.
  433. * There is no create-time per-session override tier on this wire — if one
  434. * returns (a create-options contribution), it must fold in between the two.
  435. */
  436. function targetFor(agent: Agent): WebLlmTargetRef {
  437. const installed = targets.get(agent)
  438. if (installed !== undefined) return installed
  439. const logged = agent.session.requestHeader()?.config
  440. const target: WebLlmTargetRef = {
  441. current: logged === undefined
  442. ? { provider: defaults.provider, model: defaults.model }
  443. : {
  444. provider: logged.provider,
  445. model: logged.model,
  446. ...logged.reasoningEffort === undefined
  447. ? {}
  448. : { reasoningEffort: logged.reasoningEffort },
  449. },
  450. assembled: undefined,
  451. }
  452. installAgentLlmTarget(agent.ctx, target)
  453. targets.set(agent, target)
  454. return target
  455. }
  456. /** Pre-publication setup used by both fresh and resumed Web agents. */
  457. function installTarget(agentCtx: Context): void {
  458. const agent = agentCtx.agent
  459. if (agent === undefined) throw new Error('api-proxy: agent setup has no scoped agent')
  460. targetFor(agent)
  461. }
  462. /** Send one transient frame to every connected mux consumer. */
  463. function broadcast(payload: MuxFrame): void {
  464. const envelope = frame(payload)
  465. for (const queue of muxQueues) queue.push(envelope)
  466. }
  467. // Projection change feed → session/projection push frames. The carrier
  468. // mints the wire frame (the seam package holds no wire vocabulary); the
  469. // child activates only when a projection registry is composed, and the
  470. // subscription unwinds with this gateway's fiber.
  471. ctx.inject(['sessionProjections'], (projectionCtx) => {
  472. projectionCtx.sessionProjections.onChanged((session, key, value, seq) => {
  473. broadcast({ type: 'session/projection', sessionId: session.id, key, value, seq })
  474. })
  475. })
  476. /**
  477. * Per-session inbox occurrence mirror serving the mux-open queue snapshot
  478. * (the same refresh-recovery baseline as pending questions). Each terminal
  479. * inbox event retires one matching occurrence, so repeated sends of the same
  480. * identified message remain visible until every occurrence is claimed.
  481. */
  482. const queuedMirror = new Map<SessionId, { message: UserMessage; steering: boolean }[]>()
  483. ctx.effect(() => {
  484. const retire = (agent: Agent, id: MessageId, placement?: InboxPlacement): void => {
  485. const entries = queuedMirror.get(agent.id)
  486. if (entries === undefined) return
  487. const index = entries.findIndex(entry =>
  488. entry.message.id === id
  489. && (placement === undefined || entry.steering === (placement === 'steering')))
  490. if (index !== -1) entries.splice(index, 1)
  491. if (entries.length === 0) queuedMirror.delete(agent.id)
  492. }
  493. const disposers = [
  494. ctx.on('agent/inbox/enqueue', (agent: Agent, message: UserMessage, placement) => {
  495. let entries = queuedMirror.get(agent.id)
  496. if (entries === undefined) {
  497. entries = []
  498. queuedMirror.set(agent.id, entries)
  499. }
  500. const steering = placement === 'steering'
  501. entries.push({ message, steering })
  502. broadcast({
  503. type: 'session/queued',
  504. sessionId: agent.id,
  505. message,
  506. steering,
  507. })
  508. }),
  509. ctx.on('agent/inbox/dequeue', (agent: Agent, message: UserMessage, placement) => {
  510. retire(agent, message.id, placement)
  511. }),
  512. ctx.on('agent/inbox/discard', (agent: Agent, messages: UserMessage[]) => {
  513. for (const message of messages) retire(agent, message.id)
  514. }),
  515. ctx.on('session/disposed', (session: Session) => {
  516. queuedMirror.delete(session.id)
  517. }),
  518. ]
  519. return () => { for (const dispose of disposers) dispose() }
  520. }, 'api-proxy: queued mirror')
  521. /** Remove a wait before settling it: synchronous deletion makes the first claimant win. */
  522. function claimQuestion(pending: PendingQuestion, outcome: 'answered' | 'cancelled'): void {
  523. pendingQuestions.delete(pending.rpcId)
  524. if (pending.signal !== undefined && pending.onAbort !== undefined) {
  525. pending.signal.removeEventListener('abort', pending.onAbort)
  526. }
  527. broadcast({
  528. type: 'question/resolved', sessionId: pending.sessionId,
  529. questionRpcId: pending.rpcId, outcome,
  530. })
  531. }
  532. const disposeProvider = ctx.userInteraction.registerProvider({
  533. ask(request: AskUserQuestionRequest): Promise<AskUserQuestionAnswer> {
  534. const sessionId = request.agent?.id
  535. if (sessionId === undefined) {
  536. return Promise.reject(new UserInteractionError(
  537. 'web user interaction requires an agent-owned session', 'ASK_MISSING_AGENT'))
  538. }
  539. return new Promise<AskUserQuestionAnswer>((resolve, reject) => {
  540. const rpcId = RpcId(randomUUID())
  541. const pending: PendingQuestion = {
  542. rpcId, sessionId, questions: request.questions, resolve, reject,
  543. ...(request.signal === undefined ? {} : { signal: request.signal }),
  544. }
  545. const onAbort = (): void => {
  546. claimQuestion(pending, 'cancelled')
  547. reject(new UserInteractionError(
  548. 'ask_user_question was aborted before the user answered', 'ASK_ABORTED'))
  549. }
  550. pending.onAbort = onAbort
  551. pendingQuestions.set(rpcId, pending)
  552. request.signal?.addEventListener('abort', onAbort, { once: true })
  553. const envelope: RpcRequest<MuxFrame> = {
  554. rpcId,
  555. payload: { type: 'question/requested', sessionId, questions: request.questions },
  556. }
  557. for (const queue of muxQueues) queue.push(envelope)
  558. })
  559. },
  560. })
  561. ctx.effect(() => () => {
  562. disposeProvider()
  563. for (const pending of [...pendingQuestions.values()]) {
  564. claimQuestion(pending, 'cancelled')
  565. pending.reject(new UserInteractionError(
  566. 'web user-interaction provider was disposed', 'ASK_ABORTED'))
  567. }
  568. }, 'api-proxy: user-interaction provider')
  569. // --- Approval pending registry ------------------------------------------
  570. // The proxy is the approval channel for every agent this host owns: an ask
  571. // through `ctx.approval` becomes an answerable server-request on the mux
  572. // stream (stable rpcId), settled by POST /api/respond. The entry survives
  573. // client disconnects — mux-open replays still-pending requested frames with
  574. // the same rpcId (the refresh-recovery baseline) — and withdraws on the
  575. // ask's own abort signal (turn cancel), pushing `cancelled` to subscribers.
  576. if (ctx.get('approval') !== undefined) {
  577. // Teardown parity with the question provider above: a gateway disposed
  578. // while approvals are pending settles every entry as 'cancelled' (the
  579. // service's fail-closed vocabulary), so no ask promise dangles past the
  580. // proxy's lifetime and subscribers see the withdrawal.
  581. ctx.effect(() => () => {
  582. for (const pending of [...pendingApprovals.values()]) pending.resolve('cancelled')
  583. }, 'api-proxy: approval registry teardown')
  584. ctx.on('approval/request', (req, next) => {
  585. // Dispatch rides a microtask behind the service's own signal check: an
  586. // abort landing in that window would register the abort listener AFTER
  587. // the signal fired — never invoked, entry pending forever, zombie frame
  588. // on every mux replay. Settle synchronously instead of publishing.
  589. if (req.signal?.aborted === true) return Promise.resolve<ApprovalOutcome>('cancelled')
  590. // The audit pair `approval/asked` is already appended by the service
  591. // before dispatch, but dispatch rides a microtask: parallel tool calls
  592. // can append several asked events before any answerer runs. THIS
  593. // request's event is therefore the newest asked event that is still
  594. // undecided, unclaimed by another pending entry, and — when the ask
  595. // names a call — carries the same callId.
  596. const events = req.agent.session.events
  597. const claimed = new Set<ApprovalRequestId>()
  598. for (const entry of pendingApprovals.values()) claimed.add(entry.approvalId)
  599. const decided = new Set<ApprovalRequestId>()
  600. let approvalId: ApprovalRequestId | undefined
  601. for (let i = events.length - 1; i >= 0; i -= 1) {
  602. const event = events[i] as SessionEvent
  603. if (event.type === 'approval/decided') {
  604. decided.add(event.data.id)
  605. } else if (event.type === 'approval/asked') {
  606. if (decided.has(event.data.id) || claimed.has(event.data.id)) continue
  607. // Symmetric pairing: a callId-bearing ask only takes its own call's
  608. // record, and a callId-less ask only takes a callId-less record —
  609. // so neither shape can steal the other's audit id under parallel
  610. // asks. (Today every producer — the tool executor — passes callId;
  611. // the callId-less arm guards any future non-tool asker.)
  612. if ((req.callId ?? null) !== (event.data.callId ?? null)) continue
  613. approvalId = event.data.id
  614. break
  615. }
  616. }
  617. // No asked event means the request bypassed the service's audit path —
  618. // not this channel's question; delegate to the fail-closed default.
  619. if (approvalId === undefined) return next()
  620. const id = approvalId
  621. return new Promise<ApprovalOutcome>((resolve) => {
  622. const settle = (outcome: ApprovalOutcome): void => {
  623. /* v8 ignore next 3 -- defensive double-settle guard: respond() routes
  624. through the pending table (a settled id is not-pending before it can
  625. re-settle) and the first settle removes the abort listener, so no
  626. reachable path settles twice; kept against future settle callers. */
  627. if (!pendingApprovals.delete(pending.rpcId)) return
  628. req.signal?.removeEventListener('abort', onAbort)
  629. broadcast({ type: 'approval/resolved', sessionId: pending.sessionId, approvalId: id, outcome })
  630. // A cancelled ask was already settled by the service's own signal
  631. // race, which discards this late resolution; resolving is a no-op
  632. // there and keeps this promise from dangling forever.
  633. resolve(outcome)
  634. }
  635. const onAbort = (): void => { settle('cancelled') }
  636. const pending: PendingApproval = {
  637. rpcId: RpcId(randomUUID()),
  638. sessionId: req.agent.session.id,
  639. approvalId: id,
  640. toolName: req.toolName,
  641. ...req.callId === undefined ? {} : { callId: req.callId },
  642. ...req.reason === undefined ? {} : { reason: req.reason },
  643. resolve: settle,
  644. }
  645. pendingApprovals.set(pending.rpcId, pending)
  646. req.signal?.addEventListener('abort', onAbort, { once: true })
  647. const envelope = requestedFrame(pending)
  648. for (const queue of muxQueues) queue.push(envelope)
  649. })
  650. })
  651. }
  652. /**
  653. * Gate the cold path on the store: an id absent from it, or naming a legacy
  654. * log without a cwd (pre-release stance: not served, no compatibility), is
  655. * not-found before any resume is attempted. With the gate passed, a later
  656. * resume failure is genuinely internal. No persistence configured skips the
  657. * gate — resume itself then fails loud with its own diagnostic.
  658. */
  659. async function assertServable(sessionId: SessionId): Promise<void> {
  660. const persistence = ctx.get('sessionPersistence')
  661. if (persistence === undefined) return
  662. const meta = (await persistence.list()).find(m => m.id === sessionId)
  663. if (meta === undefined || meta.cwd === undefined) throw new SessionNotFound(`session "${sessionId}" not found`)
  664. }
  665. async function agentFor(sessionId: SessionId): Promise<{ agent: Agent } | { error: RpcError }> {
  666. const live = ctx.agents.get(sessionId)
  667. if (live !== undefined) return { agent: live }
  668. let resume = resumes.get(sessionId)
  669. if (resume === undefined) {
  670. resume = (async () => {
  671. try {
  672. await assertServable(sessionId)
  673. const handle = await ctx.agents.resume({
  674. resumeSessionId: sessionId,
  675. agentOptions,
  676. setup: installTarget,
  677. })
  678. return handle.agent
  679. } finally {
  680. resumes.delete(sessionId)
  681. }
  682. })()
  683. resumes.set(sessionId, resume)
  684. }
  685. try {
  686. return { agent: await resume }
  687. } catch (error: unknown) {
  688. if (error instanceof SessionNotFound) {
  689. return { error: { code: 'session-not-found', message: error.message, details: { sessionId } } }
  690. }
  691. // The internal details slot is contractually {}; the reason rides the message.
  692. return { error: { code: 'internal', message: `resume failed for session "${sessionId}": ${String(error)}`, details: {} } }
  693. }
  694. }
  695. /** Resolve one requested identity to a live agent, creating or resuming it once. */
  696. async function ensureSession(sessionId: SessionId, cwd: string, checkPersistedIdentity: boolean): Promise<Agent> {
  697. let creation = sessionCreations.get(sessionId)
  698. if (creation === undefined) {
  699. creation = (async () => {
  700. const live = ctx.agents.get(sessionId)
  701. if (live !== undefined) return live
  702. const persistence = checkPersistedIdentity ? ctx.get('sessionPersistence') : undefined
  703. const stored = persistence === undefined
  704. ? undefined
  705. : (await persistence.list()).find(header => header.id === sessionId)
  706. if (stored !== undefined) {
  707. if (stored.cwd !== cwd) {
  708. throw new SessionCwdConflict(sessionId, cwd, stored.cwd)
  709. }
  710. return (await ctx.agents.resume({
  711. resumeSessionId: sessionId,
  712. agentOptions,
  713. setup: installTarget,
  714. })).agent
  715. }
  716. try {
  717. await mkdir(cwd, { recursive: true })
  718. } catch (error: unknown) {
  719. throw new Error(`failed to ensure project directory "${cwd}": ${String(error)}`, { cause: error })
  720. }
  721. return (await ctx.agents.create({
  722. sessionId,
  723. agentOptions,
  724. meta: { cwd },
  725. setup: installTarget,
  726. })).agent
  727. })().catch((error: unknown) => {
  728. // Another Host entry path may have published the same identity while
  729. // this operation crossed an asynchronous persistence/filesystem step.
  730. const live = ctx.agents.get(sessionId)
  731. if (live !== undefined) return live
  732. throw error
  733. }).finally(() => {
  734. sessionCreations.delete(sessionId)
  735. })
  736. sessionCreations.set(sessionId, creation)
  737. }
  738. const agent = await creation
  739. if (agent.session.header.cwd !== cwd) {
  740. throw new SessionCwdConflict(sessionId, cwd, agent.session.header.cwd)
  741. }
  742. return agent
  743. }
  744. /** Resolve or create one path while holding the Host's workspace-create chain. */
  745. function ensureWorkspace(
  746. path: string,
  747. title: string | undefined,
  748. rejectExistingName = false,
  749. createDirectory = false,
  750. ): Promise<{ workspace: Workspace; created: boolean }> {
  751. const operation = workspaceCreationChain.then(async () => {
  752. if (rejectExistingName && title !== undefined
  753. && ctx.workspace.list().some(workspace => workspace.title === title)) {
  754. throw new WorkspaceNameConflictError(title)
  755. }
  756. if (createDirectory) {
  757. try {
  758. await mkdir(path, { recursive: true })
  759. } catch (error: unknown) {
  760. throw new WorkspaceDirectoryCreationError(
  761. `failed to create workspace directory "${path}": ${String(error)}`,
  762. )
  763. }
  764. }
  765. const existing = await ctx.workspace.resolveByPath(path)
  766. if (existing !== undefined) return { workspace: existing, created: false }
  767. return { workspace: await ctx.workspace.create(path, title), created: true }
  768. })
  769. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  770. return operation
  771. }
  772. /** Resolve the goal service; absent = the deployment did not compose @deepseek-ai/dsh-goal. */
  773. function goalService(): NonNullable<ReturnType<typeof ctx.get<'goals'>>> | { error: RpcError } {
  774. const goals = ctx.get('goals')
  775. if (goals === undefined) {
  776. return { error: { code: 'internal', message: 'goal service is absent: this deployment does not mount @deepseek-ai/dsh-goal in its composition (cordis.yml or explicit assembly)', details: {} } }
  777. }
  778. return goals
  779. }
  780. /** Map one goal-domain rejection to the wire error (stable GoalError codes ride in details). */
  781. function goalError(request: RpcRequest<unknown>, error: unknown): RpcResponse<never> {
  782. const details = error instanceof GoalError ? { goalCode: error.code } : {}
  783. return err(request, { code: 'internal', message: String(error), details })
  784. }
  785. /** Resolve a session's agent, apply one goal mutation, and acknowledge with the new CAS ref. */
  786. async function mutateGoal(
  787. request: RpcRequest<{ sessionId: SessionId }>,
  788. mutation: (goals: NonNullable<ReturnType<typeof ctx.get<'goals'>>>, agent: Agent) => CoreGoalRef,
  789. ): Promise<RpcResponse<{ ref: GoalRef }>> {
  790. const goals = goalService()
  791. if ('error' in goals) return err(request, goals.error)
  792. const found = await agentFor(request.payload.sessionId)
  793. if ('error' in found) return err(request, found.error)
  794. try {
  795. const ref = mutation(goals, found.agent)
  796. return ok(request, { ref: { id: ref.id, revision: ref.revision } })
  797. } catch (error: unknown) {
  798. return goalError(request, error)
  799. }
  800. }
  801. return {
  802. sessions: {
  803. // Attached sessions summarize from memory; persisted-but-unattached (cold)
  804. // sessions merge in from the persistence store so history survives restarts.
  805. // Legacy logs without a cwd (pre-project stance) are not served — every
  806. // session now records its project at create time.
  807. async list(request) {
  808. const items = ctx.sessions.list().map((session) => {
  809. const agent = ctx.agents.get(session.id)
  810. const projections = listProjectionsFor(ctx, session.header, session)
  811. return {
  812. ...summarize(session, agent?.status === 'running'),
  813. ...projections === undefined ? {} : { projections },
  814. }
  815. })
  816. const attached = new Set(items.map(item => item.sessionId))
  817. const persistence = ctx.get('sessionPersistence')
  818. if (persistence !== undefined) {
  819. const cold = (await persistence.list()).filter(meta => !attached.has(meta.id) && meta.cwd !== undefined)
  820. items.push(...await Promise.all(cold.map(async (meta) => {
  821. // Cold rows read the persisted projection cache only — never a
  822. // log load; a session without a cache row simply has no column.
  823. const projections = listProjectionsFor(ctx, meta, undefined)
  824. return {
  825. ...await summarizeCold(persistence, meta),
  826. ...projections === undefined ? {} : { projections },
  827. }
  828. })))
  829. }
  830. items.sort((a, b) => b.updatedAt - a.updatedAt)
  831. return ok(request, { items })
  832. },
  833. async create(request) {
  834. const sessionId = request.payload.sessionId ?? `session-${randomUUID()}` as SessionId
  835. let workspace: Workspace | undefined
  836. if (request.payload.workspaceId !== undefined) {
  837. workspace = ctx.workspace.get(brandWorkspaceId(request.payload.workspaceId))
  838. if (workspace === undefined) {
  839. return err(request, {
  840. code: 'workspace-not-found',
  841. message: `workspace "${request.payload.workspaceId}" not found`,
  842. details: { workspaceId: request.payload.workspaceId },
  843. })
  844. }
  845. }
  846. const cwd = workspace?.path ?? request.payload.cwd ?? defaults.cwd
  847. try {
  848. await ensureSession(sessionId, cwd, request.payload.sessionId !== undefined)
  849. } catch (error: unknown) {
  850. if (error instanceof SessionCwdConflict) {
  851. return err(request, {
  852. code: 'session-conflict',
  853. message: error.message,
  854. details: {
  855. sessionId: error.sessionId,
  856. requestedCwd: error.requestedCwd,
  857. ...error.existingCwd === undefined ? {} : { existingCwd: error.existingCwd },
  858. },
  859. })
  860. }
  861. return err(request, {
  862. code: 'internal',
  863. message: `failed to create session "${sessionId}": ${String(error)}`,
  864. details: {},
  865. })
  866. }
  867. if (workspace !== undefined) {
  868. try {
  869. await workspace.attachSession(sessionId)
  870. } catch (error: unknown) {
  871. return err(request, {
  872. code: 'workspace-attach-failed',
  873. message: `session "${sessionId}" was created but could not attach to workspace "${workspace.id}": ${String(error)}`,
  874. details: { sessionId, workspaceId: workspace.id },
  875. })
  876. }
  877. }
  878. return ok(request, { sessionId })
  879. },
  880. async history(request) {
  881. const { sessionId, beforeSeq, maxMessages } = request.payload
  882. const found = await agentFor(sessionId)
  883. if ('error' in found) return err(request, found.error)
  884. // Everything below the resume above is synchronous: the page slice,
  885. // the seq read, and the projection walk see one un-torn session state.
  886. const page = paginate(found.agent.session.events, beforeSeq, maxMessages ?? DEFAULT_MAX_MESSAGES)
  887. // Views are computed against the registry at pagination time; result
  888. // pairing scans within the page only (message-boundary pagination keeps
  889. // a call and its result on one page — a cross-page miss soft-falls).
  890. const entries: HistoryEntry[] = page.events.map((event) => {
  891. const view = viewFor(ctx, event, callId => backscanArgs(page.events, callId))
  892. return { event, ...view === undefined ? {} : { view } }
  893. })
  894. // Baseline rider: tail page only — loadOlder (beforeSeq present) is
  895. // the one path that never needs a fresh projection baseline.
  896. const projections = beforeSeq === undefined ? projectionsFor(ctx, found.agent) : undefined
  897. return ok(request, {
  898. events: entries,
  899. hasMore: page.hasMore,
  900. ...projections === undefined ? {} : { projections },
  901. })
  902. },
  903. async models(request) {
  904. const { sessionId } = request.payload
  905. const found = await agentFor(sessionId)
  906. if ('error' in found) return err(request, found.error)
  907. const current = targetFor(found.agent).current
  908. const catalog = await Promise.all(ctx.llm.listProviders().map(async (provider) => {
  909. try {
  910. const advertised = await ctx.llm.listModels(provider.id)
  911. const models = [...advertised]
  912. if (
  913. provider.id === current.provider
  914. && !models.some(model => model.id === current.model)
  915. ) {
  916. models.push({
  917. provider: provider.id,
  918. id: current.model,
  919. name: current.model,
  920. })
  921. }
  922. const entries = await Promise.all(models.map(async (model) => {
  923. const resolved = await ctx.llm.resolveModelInfo(provider.id, model.id)
  924. const reasoning: ModelReasoning | undefined = resolved.reasoning === undefined
  925. ? undefined
  926. : {
  927. efforts: resolved.reasoning.efforts.map(effort => ({
  928. id: effort.id,
  929. name: effort.name,
  930. ...effort.description === undefined
  931. ? {}
  932. : { description: effort.description },
  933. })),
  934. ...resolved.reasoning.defaultEffort === undefined
  935. ? {}
  936. : { defaultEffort: resolved.reasoning.defaultEffort },
  937. }
  938. return {
  939. id: model.id,
  940. name: model.name,
  941. ...model.description === undefined ? {} : { description: model.description },
  942. ...provider.id === current.provider
  943. && model.id === current.model
  944. && !advertised.some(candidate => candidate.id === current.model)
  945. ? { unlisted: true as const }
  946. : {},
  947. ...reasoning === undefined ? {} : { reasoning },
  948. }
  949. }))
  950. const group: ModelProviderGroup = {
  951. id: provider.id,
  952. name: provider.name,
  953. models: entries,
  954. }
  955. return { kind: 'group' as const, group }
  956. } catch (error: unknown) {
  957. const failure: ModelCatalogFailure = {
  958. id: provider.id,
  959. name: provider.name,
  960. message: error instanceof Error ? error.message : String(error),
  961. }
  962. return { kind: 'failure' as const, failure }
  963. }
  964. }))
  965. const groups = catalog.flatMap(item => item.kind === 'group' ? [item.group] : [])
  966. const failures = catalog.flatMap(item => item.kind === 'failure' ? [item.failure] : [])
  967. return ok(request, {
  968. current: { ...current },
  969. groups: groups.filter(group => group.models.length > 0),
  970. failures,
  971. })
  972. },
  973. async selectModel(request) {
  974. const { sessionId, provider, model, reasoningEffort } = request.payload
  975. const found = await agentFor(sessionId)
  976. if ('error' in found) return err(request, found.error)
  977. try {
  978. const resolved = await ctx.llm.resolveCallConfig({
  979. provider,
  980. model,
  981. ...reasoningEffort === undefined
  982. ? {}
  983. : { reasoningEffort: ReasoningEffortId(reasoningEffort) },
  984. })
  985. const selected: AgentLlmTarget = {
  986. provider: resolved.provider,
  987. model: resolved.model,
  988. ...resolved.reasoningEffort === undefined
  989. ? {}
  990. : { reasoningEffort: resolved.reasoningEffort },
  991. }
  992. targetFor(found.agent).current = selected
  993. return ok(request, { selected: { ...selected } })
  994. } catch (error: unknown) {
  995. return err(request, {
  996. code: 'model-unavailable',
  997. message: error instanceof Error ? error.message : String(error),
  998. details: { provider, model },
  999. })
  1000. }
  1001. },
  1002. async prompt(request) {
  1003. const { sessionId, mode, content } = request.payload
  1004. const found = await agentFor(sessionId)
  1005. if ('error' in found) return err(request, found.error)
  1006. const agent = found.agent
  1007. // The rpcId rides MessageSource into user/message (merge declaration in api/sessions.ts; provisional correlation).
  1008. const source: MessageSource = { kind: 'user', rpcId: request.rpcId }
  1009. try {
  1010. const message: UserMessage = createUserMessage({ content, source })
  1011. if (mode === 'steer') agent.steer(message)
  1012. else agent.followup(message)
  1013. } catch (error: unknown) {
  1014. // A synchronous throw from steer/followup means disposed or invalid input; surface as agent-busy with the reason attached.
  1015. return err(request, { code: 'agent-busy', message: 'prompt rejected', details: { reason: String(error) } })
  1016. }
  1017. return ok(request, { accepted: true as const })
  1018. },
  1019. cancel(request) {
  1020. const { sessionId } = request.payload
  1021. const agent = ctx.agents.get(sessionId)
  1022. if (agent === undefined) {
  1023. return Promise.resolve(err(request, {
  1024. code: 'session-not-found',
  1025. message: `session "${sessionId}" not found (not attached)`,
  1026. details: { sessionId },
  1027. }))
  1028. }
  1029. agent.cancel({ kind: 'user' })
  1030. return Promise.resolve(ok(request, { accepted: true as const }))
  1031. },
  1032. },
  1033. workspace: {
  1034. list(request) {
  1035. return Promise.resolve(ok(request, { items: ctx.workspace.list().map(workspaceView) }))
  1036. },
  1037. // Exactly one of path/name arrives (schema refine). Existing-folder
  1038. // adoption reuses its canonical path; create-by-name rejects a name
  1039. // already present in the registry.
  1040. async create(request) {
  1041. const { payload } = request
  1042. let path: string
  1043. if (payload.name !== undefined) {
  1044. const name = payload.name.trim()
  1045. if (name === '' || name === '.' || name === '..' || /[/\\]/.test(name)) {
  1046. return err(request, {
  1047. code: 'workspace-invalid-path',
  1048. message: `workspace name must be one non-empty path segment, got "${payload.name}"`,
  1049. details: { path: payload.name },
  1050. })
  1051. }
  1052. path = join(defaults.workspaceRoot, name)
  1053. } else {
  1054. path = payload.path as string
  1055. }
  1056. try {
  1057. const name = payload.name?.trim()
  1058. const { workspace, created } = await ensureWorkspace(
  1059. path,
  1060. name,
  1061. name !== undefined,
  1062. name !== undefined,
  1063. )
  1064. return ok(request, { workspace: workspaceView(workspace), created })
  1065. } catch (error: unknown) {
  1066. if (error instanceof WorkspaceNameConflictError) {
  1067. return err(request, {
  1068. code: 'workspace-name-conflict',
  1069. message: error.message,
  1070. details: { name: error.workspaceName },
  1071. })
  1072. }
  1073. if (error instanceof WorkspaceDirectoryCreationError) {
  1074. return err(request, { code: 'internal', message: error.message, details: {} })
  1075. }
  1076. // The registry rejects a path that does not resolve to an existing
  1077. // directory (realpath ENOENT / not-a-directory) — the business
  1078. // error of the typed-path flow, surfaced as a validation failure.
  1079. return err(request, {
  1080. code: 'workspace-invalid-path',
  1081. message: `cannot create a workspace at "${path}": ${error instanceof Error ? error.message : String(error)}`,
  1082. details: { path },
  1083. })
  1084. }
  1085. },
  1086. async rename(request) {
  1087. const { payload } = request
  1088. const workspace = ctx.workspace.get(brandWorkspaceId(payload.workspaceId))
  1089. if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
  1090. const title = payload.title.trim()
  1091. // Uniqueness AND the same-title no-op both ride the create chain so
  1092. // they observe the state left by earlier queued renames — checked
  1093. // up front, a queued A→A could report success while an earlier A→B
  1094. // still lands afterwards.
  1095. const operation = workspaceCreationChain.then(async () => {
  1096. if (title === workspace.title) return
  1097. if (ctx.workspace.list().some(other => other.id !== workspace.id && other.title === title)) {
  1098. throw new WorkspaceNameConflictError(title)
  1099. }
  1100. await workspace.setTitle(title)
  1101. })
  1102. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  1103. try {
  1104. await operation
  1105. } catch (error: unknown) {
  1106. if (error instanceof WorkspaceNameConflictError) {
  1107. return err(request, {
  1108. code: 'workspace-name-conflict',
  1109. message: error.message,
  1110. details: { name: error.workspaceName },
  1111. })
  1112. }
  1113. throw error
  1114. }
  1115. return ok(request, { workspace: workspaceView(workspace) })
  1116. },
  1117. async delete(request) {
  1118. const { workspaceId } = request.payload
  1119. const operation = workspaceCreationChain.then(() =>
  1120. ctx.workspace.delete(brandWorkspaceId(workspaceId)))
  1121. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  1122. if (!await operation) return workspaceNotFound(request, workspaceId)
  1123. return ok(request, { deleted: true as const })
  1124. },
  1125. async insertSessionBefore(request) {
  1126. const { payload } = request
  1127. const workspace = ctx.workspace.get(brandWorkspaceId(payload.workspaceId))
  1128. if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
  1129. try {
  1130. await workspace.insertSessionBefore(payload.sessionId, payload.beforeSessionId)
  1131. } catch (error: unknown) {
  1132. // Only the entity's unaccounted-id rejection is the business code;
  1133. // storage/durability failures propagate as internal errors.
  1134. if (!(error instanceof WorkspaceMoveInvalidError)) throw error
  1135. return err(request, {
  1136. code: 'workspace-move-invalid',
  1137. message: error.message,
  1138. details: {
  1139. workspaceId: payload.workspaceId,
  1140. sessionId: payload.sessionId,
  1141. ...payload.beforeSessionId === undefined ? {} : { beforeSessionId: payload.beforeSessionId },
  1142. },
  1143. })
  1144. }
  1145. return ok(request, { workspace: workspaceView(workspace) })
  1146. },
  1147. },
  1148. host: {
  1149. describe(request) {
  1150. // TODO(step2): version should read apps/cli's package.json; placeholder for now.
  1151. return Promise.resolve(ok(request, {
  1152. version: '0.0.1',
  1153. // Same source as session.create's fallback: the UI's default project
  1154. // must match where an unspecified-cwd session actually lands.
  1155. cwd: defaults.cwd,
  1156. provider: defaults.provider,
  1157. model: defaults.model,
  1158. attachedSessions: ctx.agents.list().length,
  1159. }))
  1160. },
  1161. async pickDirectory(request, signal) {
  1162. const capability = ctx.directoryPicker.capability()
  1163. if (capability.kind !== 'native') {
  1164. return err(request, {
  1165. code: 'directory-picker-unavailable',
  1166. message: `host.pickDirectory needs the native capability; the composed picker serves "${capability.kind}"`,
  1167. details: { capability: capability.kind },
  1168. })
  1169. }
  1170. try {
  1171. const path = await capability.pick(signal)
  1172. return ok(request, { path })
  1173. } catch (error: unknown) {
  1174. if (signal.aborted) {
  1175. return err(request, {
  1176. code: 'cancelled',
  1177. message: 'directory picker was aborted',
  1178. details: {},
  1179. })
  1180. }
  1181. return err(request, {
  1182. code: 'internal',
  1183. message: `directory picker failed: ${error instanceof Error ? error.message : String(error)}`,
  1184. details: {},
  1185. })
  1186. }
  1187. },
  1188. async listDirectory(request, signal) {
  1189. const capability = ctx.directoryPicker.capability()
  1190. if (capability.kind !== 'browse') {
  1191. return err(request, {
  1192. code: 'directory-picker-unavailable',
  1193. message: `host.listDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
  1194. details: { capability: capability.kind },
  1195. })
  1196. }
  1197. try {
  1198. // The carrier's signal follows the caller: a disconnect or timeout
  1199. // stops the backend's directory scan instead of outliving it.
  1200. return ok(request, await capability.list(request.payload.path, signal))
  1201. } catch (error: unknown) {
  1202. // An abort is the caller's own timeout/disconnect, not a server
  1203. // failure — same code pickDirectory and command.execute report.
  1204. if (signal.aborted) {
  1205. return err(request, { code: 'cancelled', message: 'directory listing was aborted', details: {} })
  1206. }
  1207. return err(request, directoryError(error))
  1208. }
  1209. },
  1210. async createDirectory(request) {
  1211. const capability = ctx.directoryPicker.capability()
  1212. if (capability.kind !== 'browse') {
  1213. return err(request, {
  1214. code: 'directory-picker-unavailable',
  1215. message: `host.createDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
  1216. details: { capability: capability.kind },
  1217. })
  1218. }
  1219. try {
  1220. return ok(request, { path: await capability.createDirectory(request.payload.path, request.payload.name) })
  1221. } catch (error: unknown) {
  1222. return err(request, directoryError(error))
  1223. }
  1224. },
  1225. async openPath(request, signal) {
  1226. try {
  1227. const open = defaults.openPath
  1228. ?? ((path: string, openSignal: AbortSignal) => openNativePath(path, openSignal))
  1229. await open(request.payload.path, signal)
  1230. return ok(request, { opened: true as const })
  1231. } catch (error: unknown) {
  1232. if (signal.aborted) {
  1233. return err(request, {
  1234. code: 'cancelled',
  1235. message: 'path open was aborted',
  1236. details: {},
  1237. })
  1238. }
  1239. return err(request, {
  1240. code: 'internal',
  1241. message: `path open failed: ${error instanceof Error ? error.message : String(error)}`,
  1242. details: {},
  1243. })
  1244. }
  1245. },
  1246. },
  1247. commands: {
  1248. // Both methods address one session's agent (agentFor keeps its
  1249. // resume-on-miss: clients only send a sessionId for a published
  1250. // session, and resume restores an existing entity).
  1251. async list(request) {
  1252. // Missing service = the deployment omitted dsh-commands from its
  1253. // composition, not an empty catalog: fail loud instead of serving [].
  1254. const commands = ctx.get('commands')
  1255. if (commands === undefined) {
  1256. return err(request, { code: 'internal', message: 'command registry is absent: this deployment does not mount @deepseek-ai/dsh-commands in its composition (cordis.yml or explicit assembly)', details: {} })
  1257. }
  1258. const found = await agentFor(request.payload.sessionId)
  1259. if ('error' in found) return err(request, found.error)
  1260. return ok(request, { commands: commands.list(found.agent) })
  1261. },
  1262. async execute(request, signal) {
  1263. const commands = ctx.get('commands')
  1264. if (commands === undefined) {
  1265. return err(request, { code: 'internal', message: 'command registry is absent: this deployment does not mount @deepseek-ai/dsh-commands in its composition (cordis.yml or explicit assembly)', details: {} })
  1266. }
  1267. const { sessionId, line } = request.payload
  1268. const found = await agentFor(sessionId)
  1269. if ('error' in found) return err(request, found.error)
  1270. try {
  1271. // Pure admission: the executor's durable command/run + command/done
  1272. // pair (broadcast on the mux stream) carries the outcome; the
  1273. // response reports whether the line resolved to a handler, plus the
  1274. // minted pairing id so the issuing client can correlate its request
  1275. // with the flow node the lifecycle events produce.
  1276. const execution = await commands.execute(found.agent, line, signal)
  1277. return ok(request, execution === undefined
  1278. ? { matched: false }
  1279. : { matched: true, commandId: execution.commandId })
  1280. } catch (error: unknown) {
  1281. if (signal.aborted) return err(request, { code: 'cancelled', message: 'command execution was aborted', details: {} })
  1282. return err(request, { code: 'internal', message: `command failed: ${String(error)}`, details: {} })
  1283. }
  1284. },
  1285. },
  1286. goals: {
  1287. // Mutations only — the read side is the 'goal' session projection.
  1288. // Every verb resolves the session's agent (agentFor: implicit cold
  1289. // resume, the command.* precedent) and acknowledges with the new CAS
  1290. // ref; the committed goal/change event carries the whole value to every
  1291. // client through the projection frames.
  1292. async create(request) {
  1293. const { objective, maxGoalRounds } = request.payload
  1294. return mutateGoal(request, (goals, agent) => goals.create(agent, {
  1295. objective,
  1296. ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
  1297. }))
  1298. },
  1299. async edit(request) {
  1300. const { ref, objective, maxGoalRounds } = request.payload
  1301. return mutateGoal(request, (goals, agent) => goals.edit(agent, ref, {
  1302. ...(objective !== undefined ? { objective } : {}),
  1303. ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
  1304. }))
  1305. },
  1306. async pause(request) {
  1307. return mutateGoal(request, (goals, agent) => goals.pause(agent, request.payload.ref))
  1308. },
  1309. async resume(request) {
  1310. return mutateGoal(request, (goals, agent) => goals.resume(agent, request.payload.ref))
  1311. },
  1312. async complete(request) {
  1313. return mutateGoal(request, (goals, agent) => goals.complete(agent, request.payload.ref))
  1314. },
  1315. async clear(request) {
  1316. const goals = goalService()
  1317. if ('error' in goals) return err(request, goals.error)
  1318. const found = await agentFor(request.payload.sessionId)
  1319. if ('error' in found) return err(request, found.error)
  1320. try {
  1321. goals.clear(found.agent, request.payload.ref)
  1322. return ok(request, { cleared: true as const })
  1323. } catch (error: unknown) {
  1324. return goalError(request, error)
  1325. }
  1326. },
  1327. },
  1328. skills: {
  1329. // Skill lookup never touches the Agent registry: the session address
  1330. // resolves to a canonical cwd from the host-resident session header, so
  1331. // listing skills cannot create or resume an agent as a side effect.
  1332. async list(request) {
  1333. const { sessionId } = request.payload
  1334. const session = ctx.sessions.get(sessionId)
  1335. if (session === undefined) {
  1336. return err(request, {
  1337. code: 'session-not-found',
  1338. message: `session "${sessionId}" not found (not attached)`,
  1339. details: { sessionId },
  1340. })
  1341. }
  1342. if (session.header.cwd === undefined) {
  1343. // Every served session records its project at create time; a
  1344. // cwd-less header is a pre-project legacy log (not served).
  1345. return err(request, { code: 'internal', message: `session "${sessionId}" has no project cwd`, details: {} })
  1346. }
  1347. const cwd = session.header.cwd
  1348. // Same stance as the commands domain: a missing service means the
  1349. // deployment omitted dsh-skill from its composition, not an empty
  1350. // catalog. ctx.get also keeps this handler independent of the gateway
  1351. // plugin's inject list (an undeclared `ctx.skills` property read
  1352. // fails the reflect proxy).
  1353. const skillRegistry = ctx.get('skills')
  1354. if (skillRegistry === undefined) {
  1355. return err(request, { code: 'internal', message: 'skill registry is absent: this deployment does not mount @deepseek-ai/dsh-skill in its composition (cordis.yml or explicit assembly)', details: {} })
  1356. }
  1357. try {
  1358. const skills = await skillRegistry.list({ cwd })
  1359. return ok(request, {
  1360. skills: skills.map(skill => ({
  1361. name: skill.name,
  1362. description: skill.description,
  1363. ...skill.whenToUse === undefined ? {} : { whenToUse: skill.whenToUse },
  1364. })),
  1365. })
  1366. } catch (error: unknown) {
  1367. return err(request, { code: 'internal', message: `skill listing failed: ${String(error)}`, details: {} })
  1368. }
  1369. },
  1370. },
  1371. events: {
  1372. mux(_request, signal) {
  1373. const queue = new FrameQueue<RpcRequest<MuxFrame>>()
  1374. muxQueues.add(queue)
  1375. for (const session of ctx.sessions.list()) {
  1376. subscribeSession(queue, session)
  1377. }
  1378. for (const pending of pendingQuestions.values()) {
  1379. queue.push({
  1380. rpcId: pending.rpcId,
  1381. payload: {
  1382. type: 'question/requested', sessionId: pending.sessionId,
  1383. questions: pending.questions,
  1384. },
  1385. })
  1386. }
  1387. // Refresh recovery: still-pending approval questions replay with their
  1388. // stable rpcId so a reconnecting client can still answer them.
  1389. for (const pending of pendingApprovals.values()) queue.push(requestedFrame(pending))
  1390. // Queue snapshot baseline (pendingQuestions precedent): frames replayed
  1391. // in arrival order per session; a reconnecting client rebuilds its
  1392. // queue view from these alone.
  1393. for (const [sessionId, entries] of queuedMirror) {
  1394. for (const entry of entries) {
  1395. queue.push(frame({
  1396. type: 'session/queued',
  1397. sessionId,
  1398. message: entry.message,
  1399. steering: entry.steering,
  1400. }))
  1401. }
  1402. }
  1403. // Per-session open-call table for result-view pairing. Bounded by the
  1404. // per-turn call count: entries clear on turn/end; a table miss (stream
  1405. // opened mid-turn) backscans the session's in-memory events instead.
  1406. const openCalls = new Map<SessionId, Map<string, { name: string; args: unknown }>>()
  1407. const disposers = [
  1408. ctx.on('session/event', (session: Session, event: SessionEvent) => {
  1409. if (event.type === 'tool/call') {
  1410. const data = event.data as ToolCallData
  1411. try {
  1412. let table = openCalls.get(session.id)
  1413. if (table === undefined) openCalls.set(session.id, table = new Map<string, { name: string; args: unknown }>())
  1414. table.set(data.callId, { name: data.name, args: JSON.parse(data.arguments) })
  1415. } catch {
  1416. // Unparseable model arguments: leave the table unset; the result view soft-falls.
  1417. }
  1418. } else if (event.type === 'turn/end') {
  1419. openCalls.delete(session.id)
  1420. }
  1421. const view = viewFor(ctx, event, callId =>
  1422. openCalls.get(session.id)?.get(callId) ?? backscanArgs(session.events, callId))
  1423. queue.push(frame({ type: 'session/event', sessionId: session.id, event, ...view === undefined ? {} : { view } }))
  1424. }),
  1425. ctx.on('session/created', (session: Session) => {
  1426. subscribeSession(queue, session)
  1427. }),
  1428. ctx.on('session/disposed', (session: Session) => {
  1429. openCalls.delete(session.id)
  1430. }),
  1431. ]
  1432. return queue.iterate(signal, () => {
  1433. muxQueues.delete(queue)
  1434. for (const dispose of disposers) dispose()
  1435. })
  1436. },
  1437. host(_request, signal) {
  1438. const queue = new FrameQueue<RpcRequest<HostFrame>>()
  1439. const committedWorkspaceIds = new Set(
  1440. ctx.workspace.list().map(workspace => String(workspace.id)),
  1441. )
  1442. const disposers = [
  1443. ctx.on('session/created', (session: Session) => {
  1444. queue.push(frame({
  1445. type: 'host/session-added',
  1446. sessionId: session.id,
  1447. // Derived at frame time like summarize(); a just-created session
  1448. // has run no turn yet, so this is constantly true in practice.
  1449. blank: sessionBlank(session),
  1450. ...session.header.parentSession === undefined ? {} : { parentSessionId: session.header.parentSession },
  1451. // cwd rides the frame so the client list needs no refresh to group the new session.
  1452. ...session.header.cwd === undefined ? {} : { cwd: session.header.cwd },
  1453. }))
  1454. }),
  1455. ctx.on('session/disposed', (session: Session) => {
  1456. queue.push(frame({ type: 'host/session-removed', sessionId: session.id }))
  1457. }),
  1458. ctx.on('agent/status', (agent: Agent, status: AgentStatus) => {
  1459. queue.push(frame({ type: 'host/session-status', sessionId: agent.id, running: status === 'running' }))
  1460. }),
  1461. ctx.on('agent/error', (agent: Agent, _turn: number, _step: number, error: unknown) => {
  1462. queue.push(frame({ type: 'host/agent-error', sessionId: agent.id, message: errorChain(error) }))
  1463. }),
  1464. ctx.on('domain/changed', (change) => {
  1465. if (change.domain !== 'workspace') return
  1466. if (change.table === '') {
  1467. if (change.operation !== 'put') return
  1468. const state = workspaceDomainState.parse(change.value)
  1469. for (const workspaceId of state.workspaceIds) {
  1470. if (committedWorkspaceIds.has(workspaceId)) continue
  1471. const workspace = ctx.workspace.get(workspaceId)
  1472. if (workspace === undefined) {
  1473. throw new Error(`committed workspace registry references missing workspace "${workspaceId}"`)
  1474. }
  1475. committedWorkspaceIds.add(workspaceId)
  1476. queue.push(frame({ type: 'host/workspace-changed', workspace: workspaceView(workspace) }))
  1477. }
  1478. return
  1479. }
  1480. if (change.table !== 'workspaces') return
  1481. if (change.operation === 'deleted') {
  1482. if (!committedWorkspaceIds.delete(change.key)) return
  1483. queue.push(frame({
  1484. type: 'host/workspace-removed',
  1485. workspaceId: change.key as WorkspaceId,
  1486. }))
  1487. return
  1488. }
  1489. if (!committedWorkspaceIds.has(change.key)) return
  1490. // Existing-entity table writes are complete attach/touch commits.
  1491. // A new entity's first put waits for the global registry write above.
  1492. queue.push(frame({
  1493. type: 'host/workspace-changed',
  1494. workspace: changedWorkspaceView(change.key, change.value),
  1495. }))
  1496. }),
  1497. ctx.on('commands/change', () => {
  1498. queue.push(frame({ type: 'host/commands-changed' }))
  1499. }),
  1500. ]
  1501. return queue.iterate(signal, () => { for (const dispose of disposers) dispose() })
  1502. },
  1503. },
  1504. respond(message: ClientResponse): Promise<RpcReceipt> {
  1505. // Route by the echoed rpcId (the wire correlation): approvals first,
  1506. // then questions — the two registries share one id space of UUIDs.
  1507. const approval = pendingApprovals.get(message.rpcId)
  1508. if (approval !== undefined) {
  1509. if (!message.result.ok) return Promise.resolve({ accepted: false, reason: 'bad-response' })
  1510. const parsed = approvalResponsePayloadSchema.safeParse(message.result.value)
  1511. // The payload's audit correlation must match the entry the rpcId routed
  1512. // to — a mismatched answer is malformed, not merely late.
  1513. if (!parsed.success || parsed.data.approvalId !== approval.approvalId || parsed.data.sessionId !== approval.sessionId) {
  1514. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  1515. }
  1516. approval.resolve(parsed.data.outcome)
  1517. return Promise.resolve({ accepted: true })
  1518. }
  1519. const pending = pendingQuestions.get(message.rpcId)
  1520. if (pending === undefined) return Promise.resolve({ accepted: false, reason: 'not-pending' })
  1521. if (!message.result.ok) {
  1522. if (message.result.error.code !== 'cancelled') {
  1523. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  1524. }
  1525. claimQuestion(pending, 'cancelled')
  1526. pending.reject(new UserInteractionError(
  1527. 'the user cancelled ask_user_question', 'ASK_CANCELLED'))
  1528. return Promise.resolve({ accepted: true })
  1529. }
  1530. const parsed = questionResponsePayloadSchema.safeParse(message.result.value)
  1531. if (!parsed.success) {
  1532. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  1533. }
  1534. const payload: QuestionResponsePayload = {
  1535. sessionId: parsed.data.sessionId,
  1536. answer: {
  1537. answers: parsed.data.answer.answers.map(answer => ({
  1538. id: answer.id,
  1539. selected: answer.selected,
  1540. ...(answer.custom === undefined ? {} : { custom: answer.custom }),
  1541. })),
  1542. },
  1543. }
  1544. if (!matchesQuestions(payload, pending)) {
  1545. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  1546. }
  1547. claimQuestion(pending, 'answered')
  1548. pending.resolve(payload.answer)
  1549. return Promise.resolve({ accepted: true })
  1550. },
  1551. }
  1552. }