api-proxy.ts 79 KB

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