api-proxy.ts 128 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501250225032504250525062507250825092510251125122513251425152516251725182519252025212522252325242525252625272528252925302531253225332534253525362537253825392540254125422543254425452546254725482549255025512552255325542555255625572558255925602561256225632564256525662567256825692570257125722573257425752576257725782579258025812582258325842585258625872588258925902591259225932594259525962597259825992600260126022603260426052606260726082609261026112612261326142615261626172618261926202621262226232624262526262627262826292630263126322633263426352636263726382639264026412642264326442645264626472648264926502651265226532654265526562657265826592660266126622663266426652666266726682669267026712672267326742675267626772678267926802681268226832684268526862687268826892690269126922693269426952696269726982699270027012702270327042705270627072708270927102711271227132714271527162717271827192720272127222723272427252726272727282729273027312732273327342735273627372738273927402741274227432744274527462747274827492750275127522753275427552756275727582759276027612762276327642765276627672768276927702771277227732774277527762777277827792780278127822783278427852786278727882789279027912792279327942795279627972798279928002801280228032804280528062807280828092810281128122813281428152816281728182819282028212822282328242825282628272828282928302831283228332834283528362837283828392840284128422843284428452846284728482849285028512852285328542855285628572858285928602861286228632864286528662867286828692870287128722873287428752876287728782879288028812882288328842885288628872888288928902891289228932894289528962897289828992900290129022903290429052906290729082909291029112912291329142915291629172918291929202921292229232924292529262927292829292930293129322933293429352936293729382939294029412942294329442945294629472948294929502951295229532954295529562957295829592960296129622963296429652966296729682969297029712972297329742975297629772978297929802981298229832984298529862987
  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 { Agent, AgentLlmTarget, AgentLlmTargetRef, AgentStatus } from '@deepseek-ai/dsh-agent'
  11. import { createUserMessage, freezeMessage, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  12. import { errorChain } from '@deepseek-ai/dsh-llm'
  13. import type { MessageSource } from '@deepseek-ai/dsh-llm'
  14. import { isAppendSurfaceEvent, lastActivityTime } from '@deepseek-ai/dsh-session'
  15. import type { Session, SessionEvent, SessionEventMap, SessionHeader, SessionId, UserMessage } from '@deepseek-ai/dsh-session'
  16. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  17. import { SessionQueryError, type SessionSearchCursor } from '@deepseek-ai/dsh-session-query'
  18. import { SubagentError } from '@deepseek-ai/dsh-subagent'
  19. import type { SubagentListEntry as CatalogSubagentListEntry } from '@deepseek-ai/dsh-subagent'
  20. import type { Workspace, WorkspaceRecord } from '@deepseek-ai/dsh-workspace'
  21. import {
  22. workspaceDomainState, workspaceRecord, WorkspaceId as brandWorkspaceId,
  23. WorkspaceMoveInvalidError, WorkspaceUnknownSessionError,
  24. } from '@deepseek-ai/dsh-workspace'
  25. // Type-only: brings the `ctx.tools` Context merge into this program (viewFor reads presenters).
  26. import { PresetMountError, UnknownPresetError } from '@deepseek-ai/dsh-agent-presets'
  27. import type {} from '@deepseek-ai/dsh-tools'
  28. import type {
  29. ApiProxy, CredentialView, GoalRef, HistoryEntry, HostFrame, ModelCatalogFailure, ModelProviderGroup,
  30. ModelReasoning, MuxFrame, QuestionResponsePayload, SessionProjectionsBlock, SessionSearchItem,
  31. QueuedInboxItem, SessionSummary, SettingsNamespaceView, SubagentAddress, ToolEventView,
  32. WorkspaceId, WorkspaceView,
  33. } from './api/index.ts'
  34. import {
  35. SESSION_SEARCH_RESULT_LIMIT,
  36. SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS,
  37. truncateUnicodeCodePoints,
  38. } from './api/session-search.ts'
  39. // Type-only: resolves `ctx.get('sessionProjections')` to the projection registry.
  40. import type {} from '@deepseek-ai/dsh-session-projection'
  41. // Type-only: resolves `ctx.get('sessionProjectionCache')` (the cold listing column).
  42. import type {} from '@deepseek-ai/dsh-session-projection-cache'
  43. // GoalError narrows domain rejections to their stable codes at the wire boundary.
  44. import { GoalError } from '@deepseek-ai/dsh-goal'
  45. import type { GoalRef as CoreGoalRef } from '@deepseek-ai/dsh-goal'
  46. // Type-only edges: resolve `ctx.get('commands')`, the `commands/change` event, and `ctx.get('skills')`.
  47. import type {} from '@deepseek-ai/dsh-commands'
  48. import type {} from '@deepseek-ai/dsh-skill'
  49. // The settings/credentials seams: brand guards run at this wire boundary; the
  50. // service reads stay optional (`ctx.get`) so a composition without either
  51. // provider still serves every other domain.
  52. import { SettingsConflictError, settingsNamespace } from '@deepseek-ai/dsh-settings'
  53. import type { SettingsDescriptor, SettingsNamespace, SettingsPathOp } from '@deepseek-ai/dsh-settings'
  54. import { credentialRef } from '@deepseek-ai/dsh-credentials'
  55. // Value edge: the rename impl narrows the title service's validation failure; the import also resolves `ctx.get('sessionTitle')`.
  56. import { SessionTitleInvalidError } from '@deepseek-ai/dsh-session-title'
  57. import type { CallId } from '@deepseek-ai/dsh-llm/brand'
  58. import type { ApprovalOutcome, ApprovalRequestId } from '@deepseek-ai/dsh-user-approval'
  59. // Side-effect type import: resolves the `approval/request` waterfall and
  60. // `ctx.get('approval')` without a value dependency on the seam (optional composition).
  61. import type {} from '@deepseek-ai/dsh-user-approval'
  62. import { approvalResponsePayloadSchema } from './api/approvals.schema.ts'
  63. import { questionResponsePayloadSchema } from './api/questions.schema.ts'
  64. import type { ClientResponse, RpcError, RpcReceipt, RpcRequest, RpcResponse } from './api/rpc.ts'
  65. import { RpcId } from './api/rpc.ts'
  66. import type {
  67. AskUserQuestionAnswer, AskUserQuestionItem, AskUserQuestionRequest,
  68. } from '@deepseek-ai/dsh-user-interaction'
  69. import { UserInteractionError } from '@deepseek-ai/dsh-user-interaction'
  70. import { DirectoryPickerError } from '@deepseek-ai/dsh-host-directory-picker'
  71. import { openNativePath, openNativeTextFile } from './native-path-opener.ts'
  72. /** Page size when history is called without maxMessages. */
  73. const DEFAULT_MAX_MESSAGES = 50
  74. /** Non-model settings namespaces intentionally served to the Web client. */
  75. const WEB_SETTINGS_NAMESPACES = ['permission'] as const
  76. /** Provider work budget: at most 100 calls and 2,000 inspected hits. */
  77. const SESSION_SEARCH_PROVIDER_CALL_LIMIT = 100
  78. /** Bound cold-log stat fan-out and settle each started batch before cancellation returns. */
  79. const COLD_SUMMARY_BATCH_SIZE = 16
  80. /** Conversation message event types (the pagination counting unit). */
  81. const MESSAGE_TYPES = new Set(['user/message', 'assistant/message'])
  82. /** Product settings intentionally exposed beside model-provider namespaces. */
  83. const PRODUCT_SETTINGS_NAMESPACES = new Set(['ui-onboarding'])
  84. /** Read live abort state across awaits without treating it as synchronously immutable. */
  85. function isAborted(signal: AbortSignal): boolean {
  86. return signal.aborted
  87. }
  88. /**
  89. * Message-boundary pagination: count maxMessages append-origin messages
  90. * backwards from the window tail. Replacement copies never entered the
  91. * conversation a reader sees — they restate a shadowed range for the model
  92. * alone — so they consume no quota; the page stays one contiguous raw range,
  93. * which keeps a compaction's log-only provenance on the same page as its
  94. * replacement. The cut is the starting seq of the oldest message group (chunks
  95. * group via sourceEventSeqs — never cut mid-message). The tail page naturally
  96. * includes the in-progress partial.
  97. */
  98. function paginate(
  99. events: readonly SessionEvent[],
  100. beforeSeq: number | undefined,
  101. maxMessages: number,
  102. ): { events: SessionEvent[]; hasMore: boolean } {
  103. const window = beforeSeq === undefined ? [...events] : events.filter(event => event.seq < beforeSeq)
  104. let count = 0
  105. let cut = 0
  106. for (let i = window.length - 1; i >= 0; i--) {
  107. const event = window[i] as SessionEvent
  108. if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue
  109. count++
  110. const sources = (event as { sourceEventSeqs?: number[] }).sourceEventSeqs
  111. const groupStart = sources !== undefined && sources.length > 0 ? Math.min(event.seq, ...sources) : event.seq
  112. if (count >= maxMessages) {
  113. cut = groupStart
  114. break
  115. }
  116. }
  117. const page = window.filter(event => event.seq >= cut)
  118. return { events: page, hasMore: cut > 0 }
  119. }
  120. /** Wrap an ok result echoing the request's rpcId. */
  121. function ok<T>(request: RpcRequest<unknown>, value: T): RpcResponse<T> {
  122. return { rpcId: request.rpcId, result: { ok: true, value } }
  123. }
  124. /**
  125. * Build the provider/model catalog over every registered route. Shared by the
  126. * session-scoped `session.models` and host-scoped `llm.models`. Catalog
  127. * membership stays advisory: an unlisted session target remains valid for
  128. * provider dispatch, but is not injected back into the selector after its
  129. * owning catalog stops advertising it. Per-provider failures ride `failures`
  130. * without failing the sound groups; groups that advertise nothing are dropped.
  131. */
  132. async function buildModelCatalog(ctx: Context): Promise<{
  133. groups: ModelProviderGroup[]
  134. failures: ModelCatalogFailure[]
  135. }> {
  136. const catalog = await Promise.all(ctx.llm.listProviders().map(async (provider) => {
  137. try {
  138. const models = await ctx.llm.listModels(provider.id)
  139. const entries = await Promise.all(models.map(async (model) => {
  140. const resolved = await ctx.llm.resolveModelInfo(provider.id, model.id)
  141. const reasoning: ModelReasoning | undefined = resolved.reasoning === undefined
  142. ? undefined
  143. : {
  144. efforts: resolved.reasoning.efforts.map(effort => ({
  145. id: effort.id,
  146. name: effort.name,
  147. ...effort.description === undefined
  148. ? {}
  149. : { description: effort.description },
  150. })),
  151. ...resolved.reasoning.defaultEffort === undefined
  152. ? {}
  153. : { defaultEffort: resolved.reasoning.defaultEffort },
  154. }
  155. return {
  156. id: model.id,
  157. name: model.name,
  158. ...model.description === undefined ? {} : { description: model.description },
  159. ...reasoning === undefined ? {} : { reasoning },
  160. }
  161. }))
  162. const group: ModelProviderGroup = {
  163. id: provider.id,
  164. name: provider.name,
  165. models: entries,
  166. }
  167. return { kind: 'group' as const, group }
  168. } catch (error: unknown) {
  169. const failure: ModelCatalogFailure = {
  170. id: provider.id,
  171. name: provider.name,
  172. message: error instanceof Error ? error.message : String(error),
  173. }
  174. return { kind: 'failure' as const, failure }
  175. }
  176. }))
  177. return {
  178. groups: catalog.flatMap(item => item.kind === 'group' ? [item.group] : []).filter(group => group.models.length > 0),
  179. failures: catalog.flatMap(item => item.kind === 'failure' ? [item.failure] : []),
  180. }
  181. }
  182. /** Wrap an error result echoing the request's rpcId. */
  183. function err<T>(request: RpcRequest<unknown>, error: RpcError): RpcResponse<T> {
  184. return { rpcId: request.rpcId, result: { ok: false, error } }
  185. }
  186. /** Simple async queue: core callbacks push, the AsyncIterable pulls; abort/return cleans up. */
  187. class FrameQueue<F> {
  188. private buffer: F[] = []
  189. private waiter: (() => void) | undefined
  190. private done = false
  191. push(item: F): void {
  192. if (this.done) return
  193. this.buffer.push(item)
  194. this.waiter?.()
  195. }
  196. end(): void {
  197. this.done = true
  198. this.waiter?.()
  199. }
  200. async *iterate(signal: AbortSignal, cleanup: () => void): AsyncGenerator<F> {
  201. const onAbort = (): void => { this.end() }
  202. signal.addEventListener('abort', onAbort, { once: true })
  203. try {
  204. while (true) {
  205. while (this.buffer.length > 0) yield this.buffer.shift() as F
  206. if (this.done || signal.aborted) return
  207. await new Promise<void>((resolve) => { this.waiter = resolve })
  208. this.waiter = undefined
  209. }
  210. } finally {
  211. signal.removeEventListener('abort', onAbort)
  212. cleanup()
  213. }
  214. }
  215. }
  216. /**
  217. * Server-side frame mint: pure pushes get a fresh rpcId per frame (answerable
  218. * frames — approval/question requested — mint their stable id in their
  219. * pending registries instead).
  220. */
  221. function frame<F>(payload: F): RpcRequest<F> {
  222. return { rpcId: RpcId(randomUUID()), payload }
  223. }
  224. /** Queue the subscription baseline frame. */
  225. function subscribeSession(queue: FrameQueue<RpcRequest<MuxFrame>>, session: Session): void {
  226. queue.push(frame({ type: 'session/subscribed', sessionId: session.id, lastSeq: session.seq - 1 }))
  227. }
  228. /**
  229. * Whether the session's conversation has started: no turn has run yet (a
  230. * turn is one model-loop execution). Standalone plugin events — command
  231. * lifecycle records, plan/mode, titles, goals — never open a turn, so
  232. * running `/plan` or `/goal` on a fresh session keeps it blank
  233. * (list-hidden, reusable).
  234. */
  235. function sessionBlank(session: Session): boolean {
  236. return !session.events.some(event => event.type === 'turn/start')
  237. }
  238. /** Shared Session-header projection for list baselines and creation frames. */
  239. function sessionListFields(header: SessionHeader): {
  240. parentSessionId?: SessionId
  241. origin?: 'subagent'
  242. cwd?: string
  243. agentPreset?: string
  244. } {
  245. return {
  246. ...header.parentSession === undefined ? {} : { parentSessionId: header.parentSession },
  247. ...header.origin === undefined ? {} : { origin: header.origin },
  248. ...header.cwd === undefined ? {} : { cwd: header.cwd },
  249. ...header.agentPreset === undefined ? {} : { agentPreset: header.agentPreset },
  250. }
  251. }
  252. /** SessionSummary projection for attached (in-memory) sessions. */
  253. function summarize(session: Session, running: boolean): SessionSummary {
  254. return {
  255. sessionId: session.id,
  256. // Excludes end-seed: a resumed-but-untouched session
  257. // must not sort as freshly worked in.
  258. updatedAt: lastActivityTime(session.events) ?? session.header.createdAt,
  259. running,
  260. blank: sessionBlank(session),
  261. ...sessionListFields(session.header),
  262. }
  263. }
  264. /**
  265. * SessionSummary projection for cold (persisted, unattached) sessions.
  266. * updatedAt is the log file's mtime; backends without a per-session file
  267. * (locate() undefined) fall back to the header's createdAt.
  268. */
  269. async function summarizeCold(
  270. persistence: SessionPersistence,
  271. meta: SessionHeader,
  272. signal?: AbortSignal,
  273. ): Promise<SessionSummary> {
  274. signal?.throwIfAborted()
  275. let updatedAt = meta.createdAt
  276. const location = persistence.locate(meta)
  277. signal?.throwIfAborted()
  278. if (location !== undefined) {
  279. try {
  280. updatedAt = (await stat(location.path)).mtimeMs
  281. } catch {
  282. // The log vanished between list() and stat() (concurrent cleanup); createdAt stands in.
  283. }
  284. signal?.throwIfAborted()
  285. }
  286. return {
  287. sessionId: meta.id,
  288. updatedAt,
  289. running: false,
  290. // Lazy persistence keeps never-appended sessions out of list(); reading
  291. // a cold log to check for turns would defeat the index read, so a listed
  292. // cold session is served as not-blank (its log holds its conversation).
  293. blank: false,
  294. ...meta.parentSession === undefined ? {} : { parentSessionId: meta.parentSession },
  295. ...meta.origin === undefined ? {} : { origin: meta.origin },
  296. /* v8 ignore next -- the empty arm needs a cwd-less meta, but list()
  297. filters those out (legacy logs are not served); the conditional mirrors
  298. summarize() shape. */
  299. ...meta.cwd === undefined ? {} : { cwd: meta.cwd },
  300. }
  301. }
  302. /** Map a browse-primitive failure onto the wire error vocabulary (unknown throws stay internal). */
  303. function directoryError(error: unknown): RpcError {
  304. if (error instanceof DirectoryPickerError) {
  305. return { code: error.code, message: error.message, details: { path: error.path } }
  306. }
  307. return { code: 'internal', message: error instanceof Error ? error.message : String(error), details: {} }
  308. }
  309. /** Resolved Host routing and project-directory defaults consumed by the API implementation. */
  310. export interface ApiProxyDefaults {
  311. provider: string
  312. model: string
  313. /** Default project directory for new sessions whose create request carries no cwd. */
  314. cwd: string
  315. /** Parent directory for name-created workspaces. */
  316. workspaceRoot: string
  317. /** Native open-with-default-application; injectable for carrier tests. */
  318. openPath?: (path: string, signal: AbortSignal) => Promise<void>
  319. /** Native text-editor handoff; injectable for settings-document tests. */
  320. openTextFile?: (path: string, signal: AbortSignal) => Promise<void>
  321. }
  322. /** The tool/call payload fields the presenter path reads. */
  323. interface ToolCallData { callId: string; name: string; arguments: string }
  324. /**
  325. * One outstanding approval question: the stable server-request id, the frame
  326. * material replayed to late mux subscribers, and the resolver that settles the
  327. * answerer's promise back into `ctx.approval`.
  328. */
  329. interface PendingApproval {
  330. rpcId: RpcId
  331. sessionId: SessionId
  332. approvalId: ApprovalRequestId
  333. toolName: string
  334. callId?: CallId
  335. reason?: string
  336. resolve(outcome: ApprovalOutcome): void
  337. }
  338. /** Project a pending entry into its answerable mux frame (initial push and mux-open replay share it). */
  339. function requestedFrame(pending: PendingApproval): RpcRequest<MuxFrame> {
  340. return {
  341. rpcId: pending.rpcId,
  342. payload: {
  343. type: 'approval/requested',
  344. sessionId: pending.sessionId,
  345. approvalId: pending.approvalId,
  346. toolName: pending.toolName,
  347. ...pending.callId === undefined ? {} : { callId: pending.callId },
  348. ...pending.reason === undefined ? {} : { reason: pending.reason },
  349. },
  350. }
  351. }
  352. /** One host-owned question wait, addressed by the stable server-request id. */
  353. interface PendingQuestion {
  354. rpcId: RpcId
  355. sessionId: SessionId
  356. questions: AskUserQuestionItem[]
  357. resolve: (answer: AskUserQuestionAnswer) => void
  358. reject: (error: UserInteractionError) => void
  359. signal?: AbortSignal
  360. onAbort?: () => void
  361. }
  362. /** Validate one answer batch against the exact question request it resolves. */
  363. function matchesQuestions(payload: QuestionResponsePayload, pending: PendingQuestion): boolean {
  364. if (payload.sessionId !== pending.sessionId) return false
  365. const answers = payload.answer.answers
  366. if (answers.length !== pending.questions.length) return false
  367. return answers.every((answer, index) => {
  368. const question = pending.questions[index] as AskUserQuestionItem
  369. if (answer.id !== question.id) return false
  370. if (new Set(answer.selected).size !== answer.selected.length) return false
  371. const custom = answer.custom?.trim()
  372. if (custom !== undefined && custom === '') return false
  373. if (question.multiSelect !== true) {
  374. if (custom !== undefined && answer.selected.length > 0) return false
  375. if (answer.selected.length > 1) return false
  376. }
  377. const labels = new Set(question.options?.map(option => option.label) ?? [])
  378. return answer.selected.every(label => labels.has(label))
  379. })
  380. }
  381. /**
  382. * Compute the render intent for a tool/call or tool/result event through the
  383. * presenters registered at this moment; every other event type gets none. A
  384. * result's presenter needs its call's parsed args — `argsFor` supplies them
  385. * (live: the per-session call table; history: an in-page backscan), returning
  386. * undefined when the pairing is unavailable (e.g. the call fell off the page),
  387. * which soft-falls to no view. Presenter or JSON.parse throws also soft-fall:
  388. * the client's documented default (generic JSON card) covers every miss.
  389. */
  390. function viewFor(
  391. ctx: Context,
  392. event: SessionEvent,
  393. argsFor: (callId: string) => unknown,
  394. // The presenter lives with the definition, and definitions are per agent
  395. // now: a preset registers its tools into that agent's layer, leaving the
  396. // global layer empty. Looking one up without the owner finds nothing, and
  397. // every card silently degrades to the generic renderer.
  398. agent?: Agent,
  399. ): ToolEventView | undefined {
  400. try {
  401. if (event.type === 'tool/call') {
  402. const { name, arguments: raw } = event.data as ToolCallData
  403. const view = ctx.tools.get(name, agent)?.presentCall?.(JSON.parse(raw))
  404. return view === undefined ? undefined : { for: 'call', view }
  405. }
  406. if (event.type === 'tool/result') {
  407. const { message, meta } = event.data
  408. const [result] = message.content
  409. const callId = message.source.callId
  410. const call = argsFor(callId) as { name: string; args: unknown } | undefined
  411. if (call === undefined) return undefined
  412. const view = ctx.tools.get(call.name, agent)?.presentResult?.(call.args, {
  413. content: result.content,
  414. isError: result.isError === true,
  415. ...meta === undefined ? {} : { meta },
  416. })
  417. return view === undefined ? undefined : { for: 'result', view }
  418. }
  419. } catch (error: unknown) {
  420. // A throwing presenter (or unparseable arguments) must not break delivery;
  421. // the event still ships, just without a view.
  422. console.error(`api-proxy: presenter failed for ${event.type}, falling back to generic: ${String(error)}`)
  423. }
  424. return undefined
  425. }
  426. /**
  427. * Resolve a tool/result's call pairing by scanning a window of events backwards
  428. * for the matching tool/call. Used by the history path (the page is the
  429. * window — a cross-page pairing soft-falls to no view) and by live-path table
  430. * misses after a reconnect-eviction.
  431. */
  432. function backscanArgs(events: readonly SessionEvent[], callId: string): { name: string; args: unknown } | undefined {
  433. for (let i = events.length - 1; i >= 0; i--) {
  434. const event = events[i] as SessionEvent
  435. if (event.type !== 'tool/call') continue
  436. const data = event.data as ToolCallData
  437. if (data.callId !== callId) continue
  438. try {
  439. return { name: data.name, args: JSON.parse(data.arguments) }
  440. } catch {
  441. // Unparseable stored arguments: same soft-fall as a live parse failure.
  442. return undefined
  443. }
  444. }
  445. return undefined
  446. }
  447. /** Render one detached history page through the same presenter path as ordinary history. */
  448. function historyPage(
  449. ctx: Context,
  450. events: readonly SessionEvent[],
  451. beforeSeq: number | undefined,
  452. maxMessages: number | undefined,
  453. agent?: Agent,
  454. ): { events: HistoryEntry[]; hasMore: boolean } {
  455. const page = paginate(events, beforeSeq, maxMessages ?? DEFAULT_MAX_MESSAGES)
  456. return {
  457. events: page.events.map((event) => {
  458. const view = viewFor(ctx, event, callId => backscanArgs(page.events, callId), agent)
  459. return { event, ...view === undefined ? {} : { view } }
  460. }),
  461. hasMore: page.hasMore,
  462. }
  463. }
  464. /**
  465. * The projection baseline for one history tail page: the registry's
  466. * watermark-cache snapshot — one fully synchronous read (no await between the
  467. * page slice and this), so all values and `asOfSeq` form a single consistent
  468. * cut and `asOfSeq` equals the window tail event seq. The carrier holds zero
  469. * domain knowledge (each value passed its unit's own schema inside the
  470. * registry). An absent registry means the deployment has no projection seam:
  471. * the whole block is absent and clients treat every key as capability-absent.
  472. */
  473. function projectionsFor(ctx: Context, session: Session): SessionProjectionsBlock | undefined {
  474. const registry = ctx.get('sessionProjections')
  475. if (registry === undefined) return undefined
  476. return registry.snapshot(session)
  477. }
  478. /**
  479. * The projection baseline of one session.list row, fail-soft: attached
  480. * sessions cut the registry's live watermark cache; cold sessions view the
  481. * persisted projection cache's identity-checked stored rows (zero log loads
  482. * either way — the listing use case the cache exists for). The block shape
  483. * (values + asOfSeq) matches the history tail's, so a client seeds its
  484. * value store under the same higher-seq-wins rule. Any failure — and an
  485. * empty value set — yields an absent block: a listing without projections
  486. * is degraded, never broken.
  487. */
  488. function listProjectionsFor(ctx: Context, meta: SessionHeader, session: Session | undefined): SessionProjectionsBlock | undefined {
  489. try {
  490. const block = session !== undefined
  491. ? ctx.get('sessionProjections')?.snapshot(session)
  492. : ctx.get('sessionProjectionCache')?.cachedSnapshot(meta)
  493. return block !== undefined && Object.keys(block.values).length > 0 ? block : undefined
  494. } catch (error) {
  495. ctx.logger.warn(`session.list: projection column for "${meta.id}" failed (serving the row without it): ${String(error)}`)
  496. return undefined
  497. }
  498. }
  499. /** Projection baseline for a detached history tail without Agent activation. */
  500. function detachedProjectionsFor(
  501. ctx: Context,
  502. events: readonly SessionEvent[],
  503. ): SessionProjectionsBlock | undefined {
  504. const registry = ctx.get('sessionProjections')
  505. if (registry === undefined) return undefined
  506. return registry.restore({}, events, 0).snapshot
  507. }
  508. /** Map continuation admission failures without exposing provider details. */
  509. function subagentPromptError(
  510. request: RpcRequest<{ childSessionId: SessionId }>,
  511. error: unknown,
  512. signal: AbortSignal,
  513. ): RpcResponse<never> {
  514. const childSessionId = request.payload.childSessionId
  515. if (signal.aborted) {
  516. return err(request, { code: 'cancelled', message: 'subagent prompt was cancelled', details: {} })
  517. }
  518. if (error instanceof SubagentError) {
  519. switch (error.code) {
  520. case 'NOT_RESUMABLE':
  521. return err(request, {
  522. code: 'subagent-not-resumable',
  523. message: 'subagent cannot be resumed',
  524. details: { childSessionId },
  525. })
  526. case 'UNAUTHORIZED':
  527. return err(request, {
  528. code: 'subagent-unauthorized',
  529. message: 'subagent does not belong to this parent',
  530. details: { childSessionId },
  531. })
  532. case 'DRAINING':
  533. case 'ACTIVATION_CLOSING':
  534. case 'CONTINUATION_UNAVAILABLE':
  535. case 'PERSISTENCE_UNAVAILABLE':
  536. return err(request, {
  537. code: 'subagent-delivery-unavailable',
  538. message: 'subagent follow-up is temporarily unavailable',
  539. details: { childSessionId },
  540. })
  541. default:
  542. break
  543. }
  544. }
  545. return err(request, { code: 'internal', message: 'subagent prompt failed', details: {} })
  546. }
  547. /** Verify one address and mode against the complete direct-child catalog. */
  548. async function catalogChild(
  549. ctx: Context,
  550. address: SubagentAddress,
  551. signal?: AbortSignal,
  552. ): Promise<{
  553. entry?: Extract<CatalogSubagentListEntry, { kind: 'child' }>
  554. error?: RpcError
  555. }> {
  556. const { parentSessionId, childSessionId, mode } = address
  557. try {
  558. const entries = await ctx.subagents.listChildren(parentSessionId, signal)
  559. const entry = entries.find(candidate => candidate.id === childSessionId)
  560. if (entry === undefined || (entry.kind === 'child' && entry.mode !== mode)) {
  561. return {
  562. error: {
  563. code: 'subagent-not-found',
  564. message: `session "${childSessionId}" is not a ${mode} direct child of "${parentSessionId}"`,
  565. details: { parentSessionId, childSessionId },
  566. },
  567. }
  568. }
  569. if (entry.kind === 'diagnostic') {
  570. return {
  571. error: {
  572. code: 'subagent-catalog-diagnostic',
  573. message: `subagent "${childSessionId}" is ${entry.reason}`,
  574. details: { parentSessionId, childSessionId, reason: entry.reason },
  575. },
  576. }
  577. }
  578. return { entry }
  579. } catch (error: unknown) {
  580. if (signal?.aborted
  581. || (error instanceof SubagentError && error.code === 'CANCELLED')
  582. || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')) {
  583. return { error: { code: 'cancelled', message: 'subagent catalog read was cancelled', details: {} } }
  584. }
  585. if (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') {
  586. return {
  587. error: {
  588. code: 'subagent-not-found',
  589. message: `parent session "${parentSessionId}" was not found`,
  590. details: { parentSessionId, childSessionId },
  591. },
  592. }
  593. }
  594. return { error: { code: 'internal', message: 'subagent catalog read failed', details: {} } }
  595. }
  596. }
  597. /**
  598. * Thrown by the cold-resume path when the id names no servable session
  599. * (absent from the store, or a pre-project legacy log without a cwd).
  600. */
  601. class SessionNotFound extends Error {}
  602. /** Session identity whose lifecycle belongs to subagent routing, not generic Host resume. */
  603. class SubagentSessionOwnership extends Error {
  604. constructor(readonly sessionId: SessionId) {
  605. super(`session "${sessionId}" is a subagent session; use subagent delivery`)
  606. }
  607. }
  608. /**
  609. * The requested preset differs from the one this session already runs.
  610. *
  611. * A session's composition is fixed at creation: its history was produced under
  612. * that preset's tools, so adopting the identity under a different one would
  613. * replay tool calls the rebuilt agent cannot make. Naming a different preset
  614. * is therefore a caller error rather than a switch.
  615. */
  616. class AgentPresetConflict extends Error {
  617. constructor(
  618. readonly sessionId: SessionId,
  619. readonly requestedPreset: string,
  620. readonly existingPreset: string | undefined,
  621. ) {
  622. super(
  623. existingPreset === undefined
  624. ? `session "${sessionId}" records no agent preset, so it cannot be adopted under one; `
  625. + 'a deployment composing no roster records none on any session — '
  626. : `session "${sessionId}" already runs agent preset ${JSON.stringify(existingPreset)}; `
  627. + `requested ${JSON.stringify(requestedPreset)}. A session's preset is fixed at creation.`,
  628. )
  629. }
  630. }
  631. /** Requested identity already belongs to a session with another project cwd. */
  632. class SessionCwdConflict extends Error {
  633. constructor(
  634. readonly sessionId: SessionId,
  635. readonly requestedCwd: string,
  636. readonly existingCwd: string | undefined,
  637. ) {
  638. super(
  639. `session "${sessionId}" already exists with cwd ${JSON.stringify(existingCwd)}; `
  640. + `requested ${JSON.stringify(requestedCwd)}`,
  641. )
  642. }
  643. }
  644. /** Host failed before the registry could adopt a name-created directory. */
  645. class WorkspaceDirectoryCreationError extends Error {}
  646. /** An explicit Host naming operation would duplicate another Workspace title. */
  647. class WorkspaceNameConflictError extends Error {
  648. constructor(readonly workspaceName: string) {
  649. super(`workspace name '${workspaceName}' is already in use`)
  650. this.name = 'WorkspaceNameConflictError'
  651. }
  652. }
  653. /** Shared workspace-not-found error response of the workspace.* mutation rows. */
  654. function workspaceNotFound<T>(request: RpcRequest<unknown>, workspaceId: string): RpcResponse<T> {
  655. return err(request, {
  656. code: 'workspace-not-found',
  657. message: `workspace "${workspaceId}" not found`,
  658. details: { workspaceId },
  659. })
  660. }
  661. /** Wire projection of one workspace entity (the workspace.* value row). */
  662. function workspaceView(workspace: Workspace): WorkspaceView {
  663. return {
  664. workspaceId: workspace.id,
  665. path: workspace.path,
  666. title: workspace.title,
  667. sessionIds: [...workspace.sessionIds],
  668. createdAt: workspace.createdAt,
  669. updatedAt: workspace.updatedAt,
  670. }
  671. }
  672. /** Wire projection of the durable record carried by `domain/changed`. */
  673. function changedWorkspaceView(workspaceId: string, value: unknown): WorkspaceView {
  674. const record: WorkspaceRecord = workspaceRecord.parse(value)
  675. return {
  676. workspaceId: workspaceId as WorkspaceId,
  677. path: record.path,
  678. title: record.title,
  679. sessionIds: [...record.sessionIds],
  680. createdAt: record.createdAt,
  681. updatedAt: record.updatedAt,
  682. }
  683. }
  684. /**
  685. * Implement ApiProxy over a composed host context.
  686. * @param ctx - a context with the Host spine and Workspace registry mounted.
  687. * @param defaults - host routing and project-directory defaults.
  688. * @returns the ApiProxy implementation.
  689. */
  690. export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiProxy {
  691. const agentOptions = { provider: defaults.provider, model: defaults.model }
  692. type WebLlmTargetRef = AgentLlmTargetRef & { current: AgentLlmTarget }
  693. const targets = new WeakMap<Agent, WebLlmTargetRef>()
  694. /** Implicit resume of cold sessions, deduplicating concurrent calls (follows the jsonrpc sessionCreations precedent). */
  695. const resumes = new Map<SessionId, Promise<Agent>>()
  696. /** Client-chosen identity creation/resume, deduplicated across concurrent retries. */
  697. const sessionCreations = new Map<SessionId, Promise<Agent>>()
  698. /** Serializes path ownership and explicit title checks with Workspace mutations. */
  699. let workspaceCreationChain = Promise.resolve()
  700. const pendingQuestions = new Map<RpcId, PendingQuestion>()
  701. const pendingApprovals = new Map<RpcId, PendingApproval>()
  702. const muxQueues = new Set<FrameQueue<RpcRequest<MuxFrame>>>()
  703. /**
  704. * Install or return the session-local target that prompt assembly snapshots.
  705. * Seed order: latest logged request/header, else the host default routing.
  706. * There is no create-time per-session override tier on this wire — if one
  707. * returns (a create-options contribution), it must fold in between the two.
  708. */
  709. function targetFor(agent: Agent): WebLlmTargetRef {
  710. const installed = targets.get(agent)
  711. if (installed !== undefined) return installed
  712. const logged = agent.session.requestHeader()?.config
  713. const target: WebLlmTargetRef = {
  714. current: logged === undefined
  715. ? { provider: defaults.provider, model: defaults.model }
  716. : {
  717. provider: logged.provider,
  718. model: logged.model,
  719. ...logged.reasoningEffort === undefined
  720. ? {}
  721. : { reasoningEffort: logged.reasoningEffort },
  722. },
  723. assembled: undefined,
  724. }
  725. installAgentLlmTarget(agent.ctx, target)
  726. targets.set(agent, target)
  727. return target
  728. }
  729. /** Pre-publication setup used by both fresh and resumed Web agents. */
  730. function installTarget(agentCtx: Context): void {
  731. const agent = agentCtx.agent
  732. if (agent === undefined) throw new Error('api-proxy: agent setup has no scoped agent')
  733. targetFor(agent)
  734. }
  735. /**
  736. * Reject an attempt to run an existing session under a different preset.
  737. *
  738. * A caller that names no preset always adopts the session as it is, so the
  739. * common paths — reconnecting, resuming, retrying a create — are unaffected.
  740. * @param sessionId - the identity being adopted.
  741. * @param requested - the preset the request named, if any.
  742. * @param existing - the preset the session was created under, if any.
  743. * @throws when both are present and differ.
  744. */
  745. function assertPresetUnchanged(
  746. sessionId: SessionId,
  747. requested: string | undefined,
  748. existing: string | undefined,
  749. ): void {
  750. if (requested === undefined || requested === existing) return
  751. throw new AgentPresetConflict(sessionId, requested, existing)
  752. }
  753. /**
  754. * Resolve the preset an agent will be composed from, and the setup that
  755. * installs it.
  756. *
  757. * The id is resolved BEFORE the session exists because the session boundary
  758. * snapshots `meta` before asynchronous setup begins — a preset discovered
  759. * during setup could never reach the header. Mounting still happens in
  760. * setup, where a failure rolls the whole creation back rather than leaving a
  761. * published session whose capabilities are half-installed.
  762. *
  763. * A deployment with no preset roster composes nothing and every session
  764. * shares the host composition, which is the behavior before presets existed.
  765. * @param presetId - the requested preset, or `undefined` for the default.
  766. * @returns the id to record on the header (absent without a roster) and the setup callback.
  767. * @throws when the roster supplies no such preset.
  768. */
  769. async function composeAgent(presetId: string | undefined): Promise<{
  770. agentPreset?: string
  771. setup: (agentCtx: Context) => Promise<void>
  772. }> {
  773. const presets = ctx.get('agentPresets')
  774. if (presets === undefined) {
  775. return {
  776. setup: (agentCtx: Context) => {
  777. installTarget(agentCtx)
  778. return Promise.resolve()
  779. },
  780. }
  781. }
  782. const resolvedId = (await presets.resolve(presetId)).id
  783. return {
  784. agentPreset: resolvedId,
  785. setup: async (agentCtx: Context) => {
  786. installTarget(agentCtx)
  787. await presets.mount(agentCtx, resolvedId)
  788. },
  789. }
  790. }
  791. /** Send one transient frame to every connected mux consumer. */
  792. function broadcast(payload: MuxFrame): void {
  793. const envelope = frame(payload)
  794. for (const queue of muxQueues) queue.push(envelope)
  795. }
  796. // Projection change feed → session/projection push frames. The carrier
  797. // mints the wire frame (the seam package holds no wire vocabulary); the
  798. // child activates only when a projection registry is composed, and the
  799. // subscription unwinds with this gateway's fiber.
  800. ctx.inject(['sessionProjections'], (projectionCtx) => {
  801. projectionCtx.sessionProjections.onChanged((session, key, value, seq) => {
  802. broadcast({ type: 'session/projection', sessionId: session.id, key, value, seq })
  803. })
  804. })
  805. /** Project both durable inbox lists, optionally including the splice currently being emitted. */
  806. const queueItems = (
  807. agent: Agent,
  808. splice?: SessionEventMap['agent/inbox/spliced'],
  809. ): QueuedInboxItem[] => {
  810. const project = (target: 'next-turn' | 'next-step'): readonly UserMessage[] => {
  811. const messages = target === 'next-turn' ? agent.inbox.nextTurn : agent.inbox.nextStep
  812. return splice?.target === target
  813. ? messages.toSpliced(splice.start, splice.removedCount ?? 0, ...splice.inserted)
  814. : messages
  815. }
  816. return [
  817. ...project('next-turn').map(message => ({ id: message.id, placement: 'queued' as const, message })),
  818. ...project('next-step').map(message => ({
  819. id: message.id,
  820. // Only user-origin messages are steering; injected context (approval
  821. // notices, task completion, attached snapshots) is not a user action
  822. // and must not render as a pending steering bubble.
  823. placement: message.source.kind === 'user' ? 'steering' as const : 'context' as const,
  824. message,
  825. })),
  826. ]
  827. }
  828. ctx.on('session/event', (session, event) => {
  829. if (event.type !== 'agent/inbox/spliced') return
  830. const agent = ctx.agents.get(session.id)
  831. if (agent?.session !== session) return
  832. broadcast({ type: 'session/queue', sessionId: session.id, items: queueItems(agent, event.data) })
  833. })
  834. /** Remove a wait before settling it: synchronous deletion makes the first claimant win. */
  835. function claimQuestion(pending: PendingQuestion, outcome: 'answered' | 'cancelled'): void {
  836. pendingQuestions.delete(pending.rpcId)
  837. if (pending.signal !== undefined && pending.onAbort !== undefined) {
  838. pending.signal.removeEventListener('abort', pending.onAbort)
  839. }
  840. broadcast({
  841. type: 'question/resolved', sessionId: pending.sessionId,
  842. questionRpcId: pending.rpcId, outcome,
  843. })
  844. }
  845. const disposeProvider = ctx.userInteraction.registerProvider({
  846. ask(request: AskUserQuestionRequest): Promise<AskUserQuestionAnswer> {
  847. const sessionId = request.agent?.id
  848. if (sessionId === undefined) {
  849. return Promise.reject(new UserInteractionError(
  850. 'web user interaction requires an agent-owned session', 'ASK_MISSING_AGENT'))
  851. }
  852. return new Promise<AskUserQuestionAnswer>((resolve, reject) => {
  853. const rpcId = RpcId(randomUUID())
  854. const pending: PendingQuestion = {
  855. rpcId, sessionId, questions: request.questions, resolve, reject,
  856. ...(request.signal === undefined ? {} : { signal: request.signal }),
  857. }
  858. const onAbort = (): void => {
  859. claimQuestion(pending, 'cancelled')
  860. reject(new UserInteractionError(
  861. 'ask_user_question was aborted before the user answered', 'ASK_ABORTED'))
  862. }
  863. pending.onAbort = onAbort
  864. pendingQuestions.set(rpcId, pending)
  865. request.signal?.addEventListener('abort', onAbort, { once: true })
  866. const envelope: RpcRequest<MuxFrame> = {
  867. rpcId,
  868. payload: { type: 'question/requested', sessionId, questions: request.questions },
  869. }
  870. for (const queue of muxQueues) queue.push(envelope)
  871. })
  872. },
  873. })
  874. ctx.effect(() => () => {
  875. disposeProvider()
  876. for (const pending of [...pendingQuestions.values()]) {
  877. claimQuestion(pending, 'cancelled')
  878. pending.reject(new UserInteractionError(
  879. 'web user-interaction provider was disposed', 'ASK_ABORTED'))
  880. }
  881. }, 'api-proxy: user-interaction provider')
  882. // --- Approval pending registry ------------------------------------------
  883. // The proxy is the approval channel for every agent this host owns: an ask
  884. // through `ctx.approval` becomes an answerable server-request on the mux
  885. // stream (stable rpcId), settled by POST /api/respond. The entry survives
  886. // client disconnects — mux-open replays still-pending requested frames with
  887. // the same rpcId (the refresh-recovery baseline) — and withdraws on the
  888. // ask's own abort signal (turn cancel), pushing `cancelled` to subscribers.
  889. if (ctx.get('approval') !== undefined) {
  890. // Teardown parity with the question provider above: a gateway disposed
  891. // while approvals are pending settles every entry as 'cancelled' (the
  892. // service's fail-closed vocabulary), so no ask promise dangles past the
  893. // proxy's lifetime and subscribers see the withdrawal.
  894. ctx.effect(() => () => {
  895. for (const pending of [...pendingApprovals.values()]) pending.resolve('cancelled')
  896. }, 'api-proxy: approval registry teardown')
  897. ctx.on('approval/request', (req, next) => {
  898. // Dispatch rides a microtask behind the service's own signal check: an
  899. // abort landing in that window would register the abort listener AFTER
  900. // the signal fired — never invoked, entry pending forever, zombie frame
  901. // on every mux replay. Settle synchronously instead of publishing.
  902. if (req.signal?.aborted === true) return Promise.resolve<ApprovalOutcome>('cancelled')
  903. // The audit pair `approval/asked` is already appended by the service
  904. // before dispatch, but dispatch rides a microtask: parallel tool calls
  905. // can append several asked events before any answerer runs. THIS
  906. // request's event is therefore the newest asked event that is still
  907. // undecided, unclaimed by another pending entry, and — when the ask
  908. // names a call — carries the same callId.
  909. const events = req.agent.session.events
  910. const claimed = new Set<ApprovalRequestId>()
  911. for (const entry of pendingApprovals.values()) claimed.add(entry.approvalId)
  912. const decided = new Set<ApprovalRequestId>()
  913. let approvalId: ApprovalRequestId | undefined
  914. for (let i = events.length - 1; i >= 0; i -= 1) {
  915. const event = events[i] as SessionEvent
  916. if (event.type === 'approval/decided') {
  917. decided.add(event.data.id)
  918. } else if (event.type === 'approval/asked') {
  919. if (decided.has(event.data.id) || claimed.has(event.data.id)) continue
  920. // Symmetric pairing: a callId-bearing ask only takes its own call's
  921. // record, and a callId-less ask only takes a callId-less record —
  922. // so neither shape can steal the other's audit id under parallel
  923. // asks. (Today every producer — the tool executor — passes callId;
  924. // the callId-less arm guards any future non-tool asker.)
  925. if ((req.callId ?? null) !== (event.data.callId ?? null)) continue
  926. approvalId = event.data.id
  927. break
  928. }
  929. }
  930. // No asked event means the request bypassed the service's audit path —
  931. // not this channel's question; delegate to the fail-closed default.
  932. if (approvalId === undefined) return next()
  933. const id = approvalId
  934. return new Promise<ApprovalOutcome>((resolve) => {
  935. const settle = (outcome: ApprovalOutcome): void => {
  936. /* v8 ignore next 3 -- defensive double-settle guard: respond() routes
  937. through the pending table (a settled id is not-pending before it can
  938. re-settle) and the first settle removes the abort listener, so no
  939. reachable path settles twice; kept against future settle callers. */
  940. if (!pendingApprovals.delete(pending.rpcId)) return
  941. req.signal?.removeEventListener('abort', onAbort)
  942. broadcast({ type: 'approval/resolved', sessionId: pending.sessionId, approvalId: id, outcome })
  943. // A cancelled ask was already settled by the service's own signal
  944. // race, which discards this late resolution; resolving is a no-op
  945. // there and keeps this promise from dangling forever.
  946. resolve(outcome)
  947. }
  948. const onAbort = (): void => { settle('cancelled') }
  949. const pending: PendingApproval = {
  950. rpcId: RpcId(randomUUID()),
  951. sessionId: req.agent.session.id,
  952. approvalId: id,
  953. toolName: req.toolName,
  954. ...req.callId === undefined ? {} : { callId: req.callId },
  955. ...req.reason === undefined ? {} : { reason: req.reason },
  956. resolve: settle,
  957. }
  958. pendingApprovals.set(pending.rpcId, pending)
  959. req.signal?.addEventListener('abort', onAbort, { once: true })
  960. const envelope = requestedFrame(pending)
  961. for (const queue of muxQueues) queue.push(envelope)
  962. })
  963. })
  964. }
  965. /** Whether the session's own suffix carries the durable subagent discriminator. */
  966. function hasSubagentDescriptor(session: Pick<Session, 'events' | 'header'>): boolean {
  967. const events = session.events
  968. // Indexed scan from the own-suffix start: slicing copies the whole suffix
  969. // on every Agent-bound RPC, including each `session.prompt` on long
  970. // transcripts.
  971. for (let index = session.header.seedLength ?? 0; index < events.length; index += 1) {
  972. if (events[index]?.type === 'subagent/descriptor') return true
  973. }
  974. return false
  975. }
  976. /**
  977. * Generic Host interaction cannot claim a durably classified subagent or an
  978. * Agent created through its live parent. The runtime-owner arm also covers
  979. * descriptor-less child publication windows and older stored headers.
  980. */
  981. function hasSubagentOwner(
  982. session: Pick<Session, 'events' | 'header'>,
  983. agent: Agent | undefined,
  984. ): boolean {
  985. if (session.header.origin === 'subagent' || hasSubagentDescriptor(session)) return true
  986. const parentId = session.header.parentSession
  987. if (parentId === undefined || agent === undefined) return false
  988. const parent = ctx.agents.get(parentId)
  989. return parent !== undefined && ctx.agents.isOwnedBy(agent.id, parent)
  990. }
  991. /** Stable generic-Host error for an identity reserved to subagent routing. */
  992. function subagentOwnershipError(sessionId: SessionId): RpcError {
  993. return {
  994. code: 'agent-busy',
  995. message: `session "${sessionId}" is owned by subagent routing`,
  996. details: { reason: 'use subagent delivery for this child session' },
  997. }
  998. }
  999. /** Inspect one cold served session without repairing, resuming, or publishing it. */
  1000. async function inspectServable(sessionId: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  1001. const persistence = ctx.get('sessionPersistence')
  1002. if (persistence === undefined) {
  1003. throw new Error('session persistence is not configured (load a dsh-session-persistence backend)')
  1004. }
  1005. const meta = (await persistence.list()).find(m => m.id === sessionId)
  1006. if (meta === undefined || meta.cwd === undefined) throw new SessionNotFound(`session "${sessionId}" not found`)
  1007. const inspected = await persistence.inspect(sessionId)
  1008. if (inspected.meta.cwd === undefined) throw new SessionNotFound(`session "${sessionId}" not found`)
  1009. return { meta: inspected.meta, events: [...inspected.events] }
  1010. }
  1011. /**
  1012. * Resolve one live registered identity through the subagent-ownership
  1013. * fence: subagent-owned agents answer `agent-busy`, plain agents pass.
  1014. * Fences the live agent's own session rather than trusting a
  1015. * "registered ⇒ attached-store" invariant — a registered subagent whose
  1016. * session is ever absent from the attached store must still not be handed
  1017. * out through generic Host routing. `undefined` means no live agent.
  1018. */
  1019. function fencedLiveAgent(sessionId: SessionId): { agent: Agent } | { error: RpcError } | undefined {
  1020. const live = ctx.agents.get(sessionId)
  1021. if (live === undefined) return undefined
  1022. if (hasSubagentOwner(live.session, live)) return { error: subagentOwnershipError(sessionId) }
  1023. return { agent: live }
  1024. }
  1025. async function agentFor(sessionId: SessionId): Promise<{ agent: Agent } | { error: RpcError }> {
  1026. const fenced = fencedLiveAgent(sessionId)
  1027. if (fenced !== undefined) return fenced
  1028. const attached = ctx.sessions.get(sessionId)
  1029. if (attached !== undefined && hasSubagentOwner(attached, undefined)) {
  1030. return { error: subagentOwnershipError(sessionId) }
  1031. }
  1032. let resume = resumes.get(sessionId)
  1033. if (resume === undefined) {
  1034. resume = (async () => {
  1035. try {
  1036. const inspected = await inspectServable(sessionId)
  1037. if (hasSubagentOwner({ header: inspected.meta, events: inspected.events }, undefined)) {
  1038. throw new SubagentSessionOwnership(sessionId)
  1039. }
  1040. const publishedSession = ctx.sessions.get(sessionId)
  1041. const publishedAgent = ctx.agents.get(sessionId)
  1042. if (publishedSession !== undefined && hasSubagentOwner(publishedSession, publishedAgent)) {
  1043. throw new SubagentSessionOwnership(sessionId)
  1044. }
  1045. // Cold resume composes the preset the session recorded, for the
  1046. // same reason `session.create` does: its history was produced under
  1047. // that composition. Every generic entry point — prompt, models,
  1048. // commands — arrives here, so leaving it out meant a session opened
  1049. // after a restart ran on host tools and the deployment persona.
  1050. const handle = await ctx.agents.resume({
  1051. resumeSessionId: sessionId,
  1052. agentOptions,
  1053. setup: (await composeAgent(inspected.meta.agentPreset)).setup,
  1054. })
  1055. return handle.agent
  1056. } finally {
  1057. resumes.delete(sessionId)
  1058. }
  1059. })()
  1060. resumes.set(sessionId, resume)
  1061. }
  1062. try {
  1063. return { agent: await resume }
  1064. } catch (error: unknown) {
  1065. if (error instanceof SessionNotFound) {
  1066. return { error: { code: 'session-not-found', message: error.message, details: { sessionId } } }
  1067. }
  1068. if (error instanceof SubagentSessionOwnership) {
  1069. return { error: subagentOwnershipError(error.sessionId) }
  1070. }
  1071. // A concurrent publish can win the identity between the pre-resume
  1072. // re-check and `ctx.agents.resume` publication; the ID-collision
  1073. // rejection falls through here. Mirror ensureSession's `.catch` in
  1074. // full: classify a subagent-owned winner into the stable ownership
  1075. // error, and hand a clean plain-agent winner straight back.
  1076. const fenced = fencedLiveAgent(sessionId)
  1077. if (fenced !== undefined) return fenced
  1078. const attached = ctx.sessions.get(sessionId)
  1079. if (attached !== undefined && hasSubagentOwner(attached, undefined)) {
  1080. return { error: subagentOwnershipError(sessionId) }
  1081. }
  1082. // The internal details slot is contractually {}; the reason rides the message.
  1083. return { error: { code: 'internal', message: `resume failed for session "${sessionId}": ${String(error)}`, details: {} } }
  1084. }
  1085. }
  1086. type SessionReadState = {
  1087. id: SessionId
  1088. header: SessionHeader
  1089. events: SessionEvent[]
  1090. }
  1091. /** Read one stable session prefix without acquiring an Agent owner. */
  1092. async function readSessionState(sessionId: SessionId): Promise<SessionReadState> {
  1093. const attached = ctx.sessions.get(sessionId)
  1094. if (attached !== undefined) {
  1095. return {
  1096. id: attached.id,
  1097. header: attached.header,
  1098. events: [...attached.events],
  1099. }
  1100. }
  1101. const inspected = await inspectServable(sessionId)
  1102. return { id: inspected.meta.id, header: inspected.meta, events: inspected.events }
  1103. }
  1104. /** Resolve the Workspace inherited by a fork without making ordinary loose lineage grouped. */
  1105. async function forkWorkspace(source: Pick<Session, 'id' | 'header'>): Promise<Workspace | undefined> {
  1106. const workspaces = ctx.workspace.list()
  1107. const direct = workspaces.find(workspace => workspace.sessionIds.includes(source.id))
  1108. if (direct !== undefined || source.header.origin !== 'subagent') return direct
  1109. const lineage = await ctx.sessionQuery.traceSession(source.id)
  1110. for (const ancestor of lineage.ancestors) {
  1111. const workspace = workspaces.find(candidate => candidate.sessionIds.includes(ancestor.header.id))
  1112. if (workspace !== undefined) return workspace
  1113. }
  1114. return undefined
  1115. }
  1116. /** Read one transcript cut and optional projection baseline without acquiring an Agent owner. */
  1117. async function historyStateFor(
  1118. sessionId: SessionId,
  1119. includeProjections: boolean,
  1120. ): Promise<{ events: SessionEvent[]; projections?: SessionProjectionsBlock }> {
  1121. const attached = ctx.sessions.get(sessionId)
  1122. if (attached !== undefined) {
  1123. const events = [...attached.events]
  1124. const projections = includeProjections ? projectionsFor(ctx, attached) : undefined
  1125. return { events, ...projections === undefined ? {} : { projections } }
  1126. }
  1127. const inspected = await inspectServable(sessionId)
  1128. const projections = includeProjections ? detachedProjectionsFor(ctx, inspected.events) : undefined
  1129. return {
  1130. events: inspected.events,
  1131. ...projections === undefined ? {} : { projections },
  1132. }
  1133. }
  1134. /** Resolve one requested identity to a live agent, creating or resuming it once. */
  1135. async function ensureSession(
  1136. sessionId: SessionId,
  1137. cwd: string,
  1138. checkPersistedIdentity: boolean,
  1139. presetId?: string,
  1140. ): Promise<Agent> {
  1141. let creation = sessionCreations.get(sessionId)
  1142. if (creation === undefined) {
  1143. creation = (async () => {
  1144. const attached = ctx.sessions.get(sessionId)
  1145. const live = ctx.agents.get(sessionId)
  1146. if (attached !== undefined && hasSubagentOwner(attached, live)) {
  1147. throw new SubagentSessionOwnership(sessionId)
  1148. }
  1149. if (live !== undefined) return live
  1150. const persistence = checkPersistedIdentity ? ctx.get('sessionPersistence') : undefined
  1151. const stored = persistence === undefined
  1152. ? undefined
  1153. : (await persistence.list()).find(header => header.id === sessionId)
  1154. if (persistence !== undefined && stored !== undefined) {
  1155. const inspected = await persistence.inspect(sessionId)
  1156. // Ownership first: explicit-id adoption of a session-backed
  1157. // subagent must answer `agent-busy` regardless of the requested
  1158. // cwd (the api/commands.ts contract), not a cwd conflict.
  1159. if (hasSubagentOwner({ header: inspected.meta, events: inspected.events }, undefined)) {
  1160. throw new SubagentSessionOwnership(sessionId)
  1161. }
  1162. if (inspected.meta.cwd !== cwd) {
  1163. throw new SessionCwdConflict(sessionId, cwd, inspected.meta.cwd)
  1164. }
  1165. assertPresetUnchanged(sessionId, presetId, inspected.meta.agentPreset)
  1166. // The stored preset wins over anything the request names: a resumed
  1167. // session's history was produced under that composition, and
  1168. // rebuilding it differently would replay tool calls the model can no
  1169. // longer make.
  1170. return (await ctx.agents.resume({
  1171. resumeSessionId: sessionId,
  1172. agentOptions,
  1173. setup: (await composeAgent(inspected.meta.agentPreset)).setup,
  1174. })).agent
  1175. }
  1176. try {
  1177. await mkdir(cwd, { recursive: true })
  1178. } catch (error: unknown) {
  1179. throw new Error(`failed to ensure project directory "${cwd}": ${String(error)}`, { cause: error })
  1180. }
  1181. const composition = await composeAgent(presetId)
  1182. return (await ctx.agents.create({
  1183. sessionId,
  1184. agentOptions,
  1185. meta: {
  1186. cwd,
  1187. ...composition.agentPreset === undefined ? {} : { agentPreset: composition.agentPreset },
  1188. },
  1189. setup: composition.setup,
  1190. })).agent
  1191. })().catch((error: unknown) => {
  1192. // Another Host entry path may have published the same identity while
  1193. // this operation crossed an asynchronous persistence/filesystem step.
  1194. const live = ctx.agents.get(sessionId)
  1195. if (live !== undefined) {
  1196. if (hasSubagentOwner(live.session, live)) throw new SubagentSessionOwnership(sessionId)
  1197. return live
  1198. }
  1199. const attached = ctx.sessions.get(sessionId)
  1200. if (attached !== undefined && hasSubagentOwner(attached, undefined)) {
  1201. throw new SubagentSessionOwnership(sessionId)
  1202. }
  1203. throw error
  1204. }).finally(() => {
  1205. sessionCreations.delete(sessionId)
  1206. })
  1207. sessionCreations.set(sessionId, creation)
  1208. }
  1209. const agent = await creation
  1210. if (hasSubagentOwner(agent.session, agent)) throw new SubagentSessionOwnership(sessionId)
  1211. // Beside the cwd check for the same reason, and after the await so it
  1212. // covers every path that yields a live agent — freshly created, adopted
  1213. // live, resumed from disk, or recovered by the concurrent-creation catch.
  1214. assertPresetUnchanged(sessionId, presetId, agent.session.header.agentPreset)
  1215. if (agent.session.header.cwd !== cwd) {
  1216. throw new SessionCwdConflict(sessionId, cwd, agent.session.header.cwd)
  1217. }
  1218. return agent
  1219. }
  1220. /** Resolve or create one path while holding the Host's workspace-create chain. */
  1221. function ensureWorkspace(
  1222. path: string,
  1223. title: string | undefined,
  1224. rejectExistingName = false,
  1225. createDirectory = false,
  1226. ): Promise<{ workspace: Workspace; created: boolean }> {
  1227. const operation = workspaceCreationChain.then(async () => {
  1228. if (rejectExistingName && title !== undefined
  1229. && ctx.workspace.list().some(workspace => workspace.title === title)) {
  1230. throw new WorkspaceNameConflictError(title)
  1231. }
  1232. if (createDirectory) {
  1233. try {
  1234. await mkdir(path, { recursive: true })
  1235. } catch (error: unknown) {
  1236. throw new WorkspaceDirectoryCreationError(
  1237. `failed to create workspace directory "${path}": ${String(error)}`,
  1238. )
  1239. }
  1240. }
  1241. const existing = await ctx.workspace.resolveByPath(path)
  1242. if (existing !== undefined) return { workspace: existing, created: false }
  1243. return { workspace: await ctx.workspace.create(path, title), created: true }
  1244. })
  1245. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  1246. return operation
  1247. }
  1248. /**
  1249. * Build the session.list baseline shared by listing and search visibility.
  1250. * Attached sessions come from memory; servable cold sessions merge from
  1251. * persistence, and the final order is newest-first.
  1252. */
  1253. async function listVisibleSessionSummaries(signal?: AbortSignal): Promise<SessionSummary[]> {
  1254. signal?.throwIfAborted()
  1255. const items = ctx.sessions.list().map((session) => {
  1256. const agent = ctx.agents.get(session.id)
  1257. const projections = listProjectionsFor(ctx, session.header, session)
  1258. return {
  1259. ...summarize(session, agent?.status === 'running'),
  1260. ...projections === undefined ? {} : { projections },
  1261. }
  1262. })
  1263. signal?.throwIfAborted()
  1264. const attached = new Set(items.map(item => item.sessionId))
  1265. const persistence = ctx.get('sessionPersistence')
  1266. if (persistence !== undefined) {
  1267. const cold = (await persistence.list(signal))
  1268. .filter(meta => !attached.has(meta.id) && meta.cwd !== undefined)
  1269. signal?.throwIfAborted()
  1270. for (let offset = 0; offset < cold.length; offset += COLD_SUMMARY_BATCH_SIZE) {
  1271. signal?.throwIfAborted()
  1272. const batch = cold.slice(offset, offset + COLD_SUMMARY_BATCH_SIZE)
  1273. const settled = await Promise.allSettled(
  1274. batch.map(async (meta) => {
  1275. // Cold rows read the persisted projection cache only — never a
  1276. // log load; a session without a cache row simply has no column.
  1277. const projections = listProjectionsFor(ctx, meta, undefined)
  1278. return {
  1279. ...await summarizeCold(persistence, meta, signal),
  1280. ...projections === undefined ? {} : { projections },
  1281. }
  1282. }),
  1283. )
  1284. const summaries: SessionSummary[] = []
  1285. let rejected = false
  1286. let failure: unknown
  1287. for (const result of settled) {
  1288. if (result.status === 'fulfilled') {
  1289. summaries.push(result.value)
  1290. } else if (!rejected) {
  1291. rejected = true
  1292. failure = result.reason
  1293. }
  1294. }
  1295. if (rejected) throw failure
  1296. signal?.throwIfAborted()
  1297. items.push(...summaries)
  1298. }
  1299. }
  1300. items.sort((a, b) => b.updatedAt - a.updatedAt)
  1301. return items
  1302. }
  1303. /**
  1304. * Resolve the goal service THIS agent runs.
  1305. *
  1306. * The service is per session: an agent preset mounts it behind an `isolate`
  1307. * realm, which no host context resolves. Reading it from the root would
  1308. * answer "absent" for a session whose composition mounts it — so the lookup
  1309. * is keyed by the agent, and only a deployment composing it nowhere is
  1310. * genuinely absent.
  1311. */
  1312. function goalServiceFor(agent: Agent): NonNullable<ReturnType<typeof ctx.get<'goals'>>> | { error: RpcError } {
  1313. const presets = ctx.get('agentPresets')
  1314. const goals = presets?.serviceFor(agent, 'goals') ?? ctx.get('goals')
  1315. if (goals === undefined) {
  1316. return { error: { code: 'internal', message: 'goal service is absent: neither this session\'s agent preset nor the host composition mounts @deepseek-ai/dsh-goal', details: {} } }
  1317. }
  1318. return goals
  1319. }
  1320. /** Map one goal-domain rejection to the wire error (stable GoalError codes ride in details). */
  1321. function goalError(request: RpcRequest<unknown>, error: unknown): RpcResponse<never> {
  1322. const details = error instanceof GoalError ? { goalCode: error.code } : {}
  1323. return err(request, { code: 'internal', message: String(error), details })
  1324. }
  1325. /** Resolve a session's agent, apply one goal mutation, and acknowledge with the new CAS ref. */
  1326. async function mutateGoal(
  1327. request: RpcRequest<{ sessionId: SessionId }>,
  1328. mutation: (goals: NonNullable<ReturnType<typeof ctx.get<'goals'>>>, agent: Agent) => CoreGoalRef,
  1329. ): Promise<RpcResponse<{ ref: GoalRef }>> {
  1330. const found = await agentFor(request.payload.sessionId)
  1331. if ('error' in found) return err(request, found.error)
  1332. const goals = goalServiceFor(found.agent)
  1333. if ('error' in goals) return err(request, goals.error)
  1334. try {
  1335. const ref = mutation(goals, found.agent)
  1336. return ok(request, { ref: { id: ref.id, revision: ref.revision } })
  1337. } catch (error: unknown) {
  1338. return goalError(request, error)
  1339. }
  1340. }
  1341. /** Missing-service report shared by the settings domain (skills-domain stance). */
  1342. function settingsAbsent(): RpcError {
  1343. return { code: 'internal', message: 'settings service is absent: this deployment does not mount a settings provider (e.g. @deepseek-ai/dsh-settings-local) in its composition', details: {} }
  1344. }
  1345. /** Open one Host-resolved target and map native failures onto the wire vocabulary. */
  1346. async function openTarget(
  1347. request: RpcRequest<unknown>, path: string, signal: AbortSignal,
  1348. open: (path: string, signal: AbortSignal) => Promise<void>,
  1349. ): Promise<RpcResponse<{ opened: true }>> {
  1350. try {
  1351. await open(path, signal)
  1352. return ok(request, { opened: true as const })
  1353. } catch (error: unknown) {
  1354. if (signal.aborted) {
  1355. return err(request, {
  1356. code: 'cancelled',
  1357. message: 'path open was aborted',
  1358. details: {},
  1359. })
  1360. }
  1361. return err(request, {
  1362. code: 'internal',
  1363. message: `path open failed: ${error instanceof Error ? error.message : String(error)}`,
  1364. details: {},
  1365. })
  1366. }
  1367. }
  1368. /** Open one Host-resolved path with its default application. */
  1369. function openPath(
  1370. request: RpcRequest<unknown>, path: string, signal: AbortSignal,
  1371. ): Promise<RpcResponse<{ opened: true }>> {
  1372. const open = defaults.openPath
  1373. ?? ((target: string, openSignal: AbortSignal) => openNativePath(target, openSignal))
  1374. return openTarget(request, path, signal, open)
  1375. }
  1376. /** Open one Host-resolved text document in a native editor. */
  1377. function openTextFile(
  1378. request: RpcRequest<unknown>, path: string, signal: AbortSignal,
  1379. ): Promise<RpcResponse<{ opened: true }>> {
  1380. const open = defaults.openTextFile
  1381. ?? ((target: string, openSignal: AbortSignal) => openNativeTextFile(target, openSignal))
  1382. return openTarget(request, path, signal, open)
  1383. }
  1384. /** Missing-service report shared by the credentials domain. */
  1385. function credentialsAbsent(): RpcError {
  1386. return { code: 'internal', message: 'credentials service is absent: this deployment does not mount a credential provider (e.g. @deepseek-ai/dsh-credentials-local) in its composition', details: {} }
  1387. }
  1388. /** Map one redacted seam descriptor to its wire view. */
  1389. function namespaceView(descriptor: SettingsDescriptor): SettingsNamespaceView {
  1390. return {
  1391. ns: String(descriptor.ns),
  1392. schema: descriptor.schema,
  1393. value: descriptor.value,
  1394. ...descriptor.base === undefined ? {} : { base: descriptor.base },
  1395. ...descriptor.user === undefined ? {} : { user: descriptor.user },
  1396. applies: descriptor.applies,
  1397. secrets: (descriptor.secrets ?? []).map(secret => ({ path: [...secret.path], set: secret.set })),
  1398. revision: descriptor.revision,
  1399. }
  1400. }
  1401. /** Settings namespaces whose changes can invalidate the model catalog. */
  1402. function modelProviderNamespaces(): Set<string> {
  1403. return new Set(ctx.llm.listConfigurableProviders().map(entry => entry.settingsNs))
  1404. }
  1405. /**
  1406. * The settings namespaces this proxy serves: configurable model providers
  1407. * plus the small explicit Web preference and product-owned allowlists. The
  1408. * settings seam remains general; a future registration does not become
  1409. * remotely readable or writable by default.
  1410. */
  1411. function exposedNamespaces(): Set<string> {
  1412. const exposed = modelProviderNamespaces()
  1413. for (const ns of WEB_SETTINGS_NAMESPACES) exposed.add(ns)
  1414. for (const ns of PRODUCT_SETTINGS_NAMESPACES) exposed.add(ns)
  1415. return exposed
  1416. }
  1417. /** Refuse a namespace outside the explicit configuration-client boundary. */
  1418. function notExposed(request: RpcRequest<unknown>, ns: string): RpcResponse<SettingsNamespaceView> {
  1419. return err(request, {
  1420. code: 'settings-not-exposed',
  1421. message: `settings namespace "${ns}" is not exposed to configuration clients`,
  1422. details: { ns },
  1423. })
  1424. }
  1425. /**
  1426. * Run one settings write (merge or wholesale replace) and acknowledge with
  1427. * the namespace's new redacted view. A namespace outside the configuration
  1428. * boundary is refused before the seam is touched; every seam refusal —
  1429. * unknown or invalid namespace, read-only provider, schema validation,
  1430. * storage — becomes one `settings-rejected` carrying the seam's own message.
  1431. */
  1432. async function settingsWrite(
  1433. request: RpcRequest<unknown>,
  1434. ns: string,
  1435. mode: 'update' | 'replace' | 'mutate',
  1436. section: object,
  1437. expectedRevision?: number,
  1438. ): Promise<RpcResponse<SettingsNamespaceView>> {
  1439. const settings = ctx.get('settings')
  1440. if (settings === undefined) return err(request, settingsAbsent())
  1441. const rejected = (error: unknown): RpcResponse<SettingsNamespaceView> => {
  1442. // A stale writer is its own outcome, not a malformed request: the client
  1443. // must re-read and re-apply rather than treat the write as invalid.
  1444. if (error instanceof SettingsConflictError) {
  1445. return err(request, {
  1446. code: 'settings-conflict',
  1447. message: error.message,
  1448. details: { ns, expected: error.expected, actual: error.actual },
  1449. })
  1450. }
  1451. return err(request, {
  1452. code: 'settings-rejected',
  1453. message: error instanceof Error ? error.message : String(error),
  1454. details: { ns },
  1455. })
  1456. }
  1457. let branded: SettingsNamespace
  1458. try {
  1459. branded = settingsNamespace(ns)
  1460. } catch (error: unknown) {
  1461. // A malformed name is a client bug, reported as such; it could never be
  1462. // in the exposed set either, so naming the real fault costs no ground.
  1463. return rejected(error)
  1464. }
  1465. if (!exposedNamespaces().has(ns)) return notExposed(request, ns)
  1466. try {
  1467. if (mode === 'update') await settings.update(branded, section, expectedRevision)
  1468. else if (mode === 'replace') await settings.replace(branded, section, expectedRevision)
  1469. else await settings.mutate(branded, section as SettingsPathOp[], expectedRevision)
  1470. } catch (error: unknown) {
  1471. return rejected(error)
  1472. }
  1473. const descriptor = settings.describe({ redactSecrets: true }).find(candidate => candidate.ns === branded)
  1474. if (descriptor === undefined) {
  1475. // The write committed but the namespace vanished before this read: only
  1476. // a concurrent registrant disposal can produce it.
  1477. return err(request, { code: 'internal', message: `settings namespace "${ns}" was disposed after the ${mode}`, details: {} })
  1478. }
  1479. return ok(request, namespaceView(descriptor))
  1480. }
  1481. return {
  1482. sessions: {
  1483. // Attached sessions summarize from memory; persisted-but-unattached (cold)
  1484. // sessions merge in from the persistence store so history survives restarts.
  1485. // Legacy logs without a cwd (pre-project stance) are not served — every
  1486. // session now records its project at create time.
  1487. async list(request) {
  1488. return ok(request, { items: await listVisibleSessionSummaries() })
  1489. },
  1490. async search(request, signal) {
  1491. const cancelled = () => err<{ items: SessionSearchItem[]; hasMore: boolean }>(request, {
  1492. code: 'cancelled',
  1493. message: 'session search was aborted',
  1494. details: {},
  1495. })
  1496. if (isAborted(signal)) return cancelled()
  1497. const sessionQuery = ctx.get('sessionQuery')
  1498. if (sessionQuery === undefined) {
  1499. return err(request, {
  1500. code: 'internal',
  1501. message: 'session search is unavailable: this deployment does not mount @deepseek-ai/dsh-session-query',
  1502. details: {},
  1503. })
  1504. }
  1505. try {
  1506. const visible = await listVisibleSessionSummaries(signal)
  1507. if (isAborted(signal)) return cancelled()
  1508. if (visible.length === 0) return ok(request, { items: [], hasMore: false })
  1509. const visibleIds = new Set(visible.map(item => item.sessionId))
  1510. const authorized: SessionSearchItem[] = []
  1511. const acceptedIds = new Set<SessionId>()
  1512. const seenCursors = new Set<SessionSearchCursor>()
  1513. let cursor: SessionSearchCursor | undefined
  1514. let providerCallCount = 0
  1515. let providerPageLimit = SESSION_SEARCH_RESULT_LIMIT
  1516. while (authorized.length <= SESSION_SEARCH_RESULT_LIMIT) {
  1517. if (isAborted(signal)) return cancelled()
  1518. if (providerCallCount >= SESSION_SEARCH_PROVIDER_CALL_LIMIT) {
  1519. throw new Error(
  1520. `session search provider exceeded the ${SESSION_SEARCH_PROVIDER_CALL_LIMIT}-call work budget`,
  1521. )
  1522. }
  1523. providerCallCount++
  1524. const requestedCursor = cursor
  1525. const requestedPageLimit = providerPageLimit
  1526. let page
  1527. try {
  1528. page = await sessionQuery.searchSessions({
  1529. query: request.payload.query,
  1530. eventFilters: [
  1531. { kind: 'type', values: ['user/message', 'assistant/message'] },
  1532. { kind: 'surface', values: ['current'] },
  1533. ],
  1534. limit: requestedPageLimit,
  1535. ...requestedCursor === undefined ? {} : { cursor: requestedCursor },
  1536. }, { signal })
  1537. } catch (error: unknown) {
  1538. if (isAborted(signal)) return cancelled()
  1539. if (
  1540. requestedCursor === undefined
  1541. && error instanceof SessionQueryError
  1542. && error.code === 'SESSION_QUERY_INVALID_LIMIT'
  1543. && requestedPageLimit > 1
  1544. ) {
  1545. providerPageLimit = Math.max(1, Math.floor(requestedPageLimit / 2))
  1546. continue
  1547. }
  1548. if (
  1549. requestedCursor !== undefined
  1550. && error instanceof SessionQueryError
  1551. && error.code === 'SESSION_QUERY_STALE_CURSOR'
  1552. ) {
  1553. authorized.length = 0
  1554. acceptedIds.clear()
  1555. seenCursors.clear()
  1556. cursor = undefined
  1557. continue
  1558. }
  1559. throw error
  1560. }
  1561. if (isAborted(signal)) return cancelled()
  1562. const providerItemCount = page.items.length
  1563. if (providerItemCount > requestedPageLimit) {
  1564. throw new Error(
  1565. `session search provider returned ${providerItemCount} items; maximum is ${requestedPageLimit}`,
  1566. )
  1567. }
  1568. // Host visibility is the authorization boundary. Consume the
  1569. // provider's globally ranked stream rather than binding every
  1570. // visible id into one SQLite statement, then re-check complete
  1571. // provenance before emitting any snippet.
  1572. for (const hit of page.items) {
  1573. if (authorized.length > SESSION_SEARCH_RESULT_LIMIT) continue
  1574. if (
  1575. !visibleIds.has(hit.header.id)
  1576. || hit.bestMatch.sessionId !== hit.header.id
  1577. || hit.bestMatch.surface !== 'current'
  1578. || !MESSAGE_TYPES.has(hit.bestMatch.type)
  1579. || acceptedIds.has(hit.header.id)
  1580. ) continue
  1581. const snippet = truncateUnicodeCodePoints(
  1582. hit.bestMatch.snippet,
  1583. SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS,
  1584. )
  1585. acceptedIds.add(hit.header.id)
  1586. authorized.push({
  1587. sessionId: hit.header.id,
  1588. snippet,
  1589. })
  1590. }
  1591. const nextCursor = page.nextCursor
  1592. if (nextCursor !== undefined) {
  1593. if (seenCursors.has(nextCursor)) {
  1594. throw new Error('session search provider repeated a continuation cursor')
  1595. }
  1596. seenCursors.add(nextCursor)
  1597. }
  1598. if (authorized.length > SESSION_SEARCH_RESULT_LIMIT || nextCursor === undefined) break
  1599. cursor = nextCursor
  1600. }
  1601. return ok(request, {
  1602. items: authorized.slice(0, SESSION_SEARCH_RESULT_LIMIT),
  1603. hasMore: authorized.length > SESSION_SEARCH_RESULT_LIMIT,
  1604. })
  1605. } catch (error: unknown) {
  1606. if (
  1607. isAborted(signal)
  1608. || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')
  1609. ) return cancelled()
  1610. // XXX: Redact provider details before exposing this gateway beyond
  1611. // its current single-user local deployment.
  1612. return err(request, {
  1613. code: 'internal',
  1614. message: `session search failed: ${String(error)}`,
  1615. details: {},
  1616. })
  1617. }
  1618. },
  1619. async create(request) {
  1620. const sessionId = request.payload.sessionId ?? `session-${randomUUID()}` as SessionId
  1621. let workspace: Workspace | undefined
  1622. if (request.payload.workspaceId !== undefined) {
  1623. workspace = ctx.workspace.get(brandWorkspaceId(request.payload.workspaceId))
  1624. if (workspace === undefined) {
  1625. return err(request, {
  1626. code: 'workspace-not-found',
  1627. message: `workspace "${request.payload.workspaceId}" not found`,
  1628. details: { workspaceId: request.payload.workspaceId },
  1629. })
  1630. }
  1631. }
  1632. const cwd = workspace?.path ?? request.payload.cwd ?? defaults.cwd
  1633. const requestedPreset = request.payload.agentPreset
  1634. try {
  1635. await ensureSession(sessionId, cwd, request.payload.sessionId !== undefined, requestedPreset)
  1636. } catch (error: unknown) {
  1637. if (error instanceof AgentPresetConflict) {
  1638. return err(request, {
  1639. code: 'agent-preset-conflict',
  1640. message: error.message,
  1641. details: {
  1642. sessionId: error.sessionId,
  1643. requestedPreset: error.requestedPreset,
  1644. ...error.existingPreset === undefined ? {} : { existingPreset: error.existingPreset },
  1645. },
  1646. })
  1647. }
  1648. if (error instanceof UnknownPresetError) {
  1649. return err(request, {
  1650. code: 'agent-preset-not-found',
  1651. message: error.message,
  1652. details: { agentPreset: error.presetId, available: [...error.available] },
  1653. })
  1654. }
  1655. if (error instanceof PresetMountError) {
  1656. return err(request, {
  1657. code: 'agent-preset-invalid',
  1658. message: error.message,
  1659. details: { agentPreset: error.presetId, reason: error.reason },
  1660. })
  1661. }
  1662. if (error instanceof SessionCwdConflict) {
  1663. return err(request, {
  1664. code: 'session-conflict',
  1665. message: error.message,
  1666. details: {
  1667. sessionId: error.sessionId,
  1668. requestedCwd: error.requestedCwd,
  1669. ...error.existingCwd === undefined ? {} : { existingCwd: error.existingCwd },
  1670. },
  1671. })
  1672. }
  1673. if (error instanceof SubagentSessionOwnership) {
  1674. return err(request, subagentOwnershipError(error.sessionId))
  1675. }
  1676. return err(request, {
  1677. code: 'internal',
  1678. message: `failed to create session "${sessionId}": ${String(error)}`,
  1679. details: {},
  1680. })
  1681. }
  1682. if (workspace !== undefined) {
  1683. try {
  1684. await workspace.attachSession(sessionId)
  1685. } catch (error: unknown) {
  1686. return err(request, {
  1687. code: 'workspace-attach-failed',
  1688. message: `session "${sessionId}" was created but could not attach to workspace "${workspace.id}": ${String(error)}`,
  1689. details: { sessionId, workspaceId: workspace.id },
  1690. })
  1691. }
  1692. }
  1693. return ok(request, { sessionId })
  1694. },
  1695. async history(request) {
  1696. const { sessionId, beforeSeq, maxMessages } = request.payload
  1697. let state: { events: SessionEvent[]; projections?: SessionProjectionsBlock }
  1698. try {
  1699. state = await historyStateFor(sessionId, beforeSeq === undefined)
  1700. } catch (error: unknown) {
  1701. if (error instanceof SessionNotFound) {
  1702. return err(request, { code: 'session-not-found', message: error.message, details: { sessionId } })
  1703. }
  1704. return err(request, {
  1705. code: 'internal',
  1706. message: `history unavailable for session "${sessionId}": ${String(error)}`,
  1707. details: {},
  1708. })
  1709. }
  1710. // `ctx.get`, not `ctx.agents`: this is the COLD path, and a caller may
  1711. // serve history from storage with no agent registry composed at all.
  1712. // An absent registry means no live agent, which is the same answer a
  1713. // present one gives here — presenters fall back to the global layer.
  1714. const page = historyPage(ctx, state.events, beforeSeq, maxMessages, ctx.get('agents')?.get(sessionId))
  1715. return ok(request, {
  1716. events: page.events,
  1717. hasMore: page.hasMore,
  1718. ...state.projections === undefined ? {} : { projections: state.projections },
  1719. })
  1720. },
  1721. async models(request) {
  1722. const { sessionId } = request.payload
  1723. const found = await agentFor(sessionId)
  1724. if ('error' in found) return err(request, found.error)
  1725. const current = targetFor(found.agent).current
  1726. const { groups, failures } = await buildModelCatalog(ctx)
  1727. return ok(request, { current: { ...current }, groups, failures })
  1728. },
  1729. async selectModel(request) {
  1730. const { sessionId, provider, model, reasoningEffort } = request.payload
  1731. const found = await agentFor(sessionId)
  1732. if ('error' in found) return err(request, found.error)
  1733. try {
  1734. const resolved = await ctx.llm.resolveCallConfig({
  1735. provider,
  1736. model,
  1737. ...reasoningEffort === undefined
  1738. ? {}
  1739. : { reasoningEffort: ReasoningEffortId(reasoningEffort) },
  1740. })
  1741. const selected: AgentLlmTarget = {
  1742. provider: resolved.provider,
  1743. model: resolved.model,
  1744. ...resolved.reasoningEffort === undefined
  1745. ? {}
  1746. : { reasoningEffort: resolved.reasoningEffort },
  1747. }
  1748. targetFor(found.agent).current = selected
  1749. return ok(request, { selected: { ...selected } })
  1750. } catch (error: unknown) {
  1751. return err(request, {
  1752. code: 'model-unavailable',
  1753. message: error instanceof Error ? error.message : String(error),
  1754. details: { provider, model },
  1755. })
  1756. }
  1757. },
  1758. async rename(request) {
  1759. const { sessionId, title } = request.payload
  1760. const found = await agentFor(sessionId)
  1761. if ('error' in found) return err(request, found.error)
  1762. const titles = ctx.get('sessionTitle')
  1763. if (titles === undefined) {
  1764. return err(request, { code: 'internal', message: 'renaming is unavailable: this deployment mounts no session-title service', details: {} })
  1765. }
  1766. try {
  1767. const accepted = titles.rename(found.agent.session, title)
  1768. return ok(request, { title: accepted.title, seq: accepted.eventSeq })
  1769. } catch (error: unknown) {
  1770. // Only the input's fault maps to title-invalid (the message is
  1771. // product-user-visible in the rename dialog); liveness and disposal
  1772. // races are deployment trouble, not a bad title.
  1773. if (error instanceof SessionTitleInvalidError) {
  1774. return err(request, {
  1775. code: 'title-invalid',
  1776. message: error.message,
  1777. details: { sessionId },
  1778. })
  1779. }
  1780. return err(request, {
  1781. code: 'internal',
  1782. message: `failed to rename session "${sessionId}": ${String(error)}`,
  1783. details: {},
  1784. })
  1785. }
  1786. },
  1787. async fork(request) {
  1788. const { sessionId, atSeq } = request.payload
  1789. let source: SessionReadState
  1790. try {
  1791. source = await readSessionState(sessionId)
  1792. } catch (error: unknown) {
  1793. if (error instanceof SessionNotFound) {
  1794. return err(request, { code: 'session-not-found', message: error.message, details: { sessionId } })
  1795. }
  1796. return err(request, {
  1797. code: 'internal',
  1798. message: `fork source unavailable for session "${sessionId}": ${String(error)}`,
  1799. details: {},
  1800. })
  1801. }
  1802. const events = source.events
  1803. // An in-log anchor belongs to the turn containing it and must never
  1804. // clip backward to an earlier completed turn. Omitted and past-end
  1805. // anchors retain the last-completed-turn shortcut.
  1806. const lastSeq = events.at(-1)?.seq ?? -1
  1807. const anchoredBoundary = atSeq === undefined
  1808. ? undefined
  1809. : events.find(e => e.type === 'turn/end' && e.seq >= atSeq)
  1810. const boundary = anchoredBoundary
  1811. ?? (atSeq === undefined || atSeq > lastSeq
  1812. ? events.findLast(e => e.type === 'turn/end')
  1813. : undefined)
  1814. if (boundary === undefined) {
  1815. return err(request, {
  1816. code: 'fork-unavailable',
  1817. message: atSeq !== undefined && atSeq <= lastSeq
  1818. ? `session "${sessionId}" has not completed the turn containing event ${String(atSeq)}`
  1819. : `session "${sessionId}" has no completed turn to fork from`,
  1820. details: { sessionId },
  1821. })
  1822. }
  1823. // Extend the cut through trailing out-of-band appends (session/title,
  1824. // injections) up to the next turn/start: they are standalone events, so
  1825. // the seed stays balanced, and the child inherits a title generated
  1826. // right after the boundary turn.
  1827. let cut = boundary.seq + 1
  1828. while (cut < events.length && events[cut]?.type !== 'turn/start') cut++
  1829. let workspace: Workspace | undefined
  1830. try {
  1831. workspace = await forkWorkspace(source)
  1832. } catch (error: unknown) {
  1833. return err(request, {
  1834. code: 'internal',
  1835. message: `failed to resolve fork workspace for session "${sessionId}": ${String(error)}`,
  1836. details: {},
  1837. })
  1838. }
  1839. const childId = `session-${randomUUID()}` as SessionId
  1840. // The child inherits the parent's composition for the same reason a
  1841. // resumed session keeps its own: the seeded history was produced under
  1842. // those tools, and composing anything else would strand the tool calls
  1843. // it already carries. Now that no model-facing row sits in the host
  1844. // plane, composing nothing would leave the child with no tools at all.
  1845. const forkComposition = await composeAgent(source.header.agentPreset)
  1846. try {
  1847. await ctx.agents.create({
  1848. sessionId: childId,
  1849. seed: events.slice(0, cut),
  1850. meta: {
  1851. ...source.header.cwd === undefined ? {} : { cwd: source.header.cwd },
  1852. parentSession: source.id,
  1853. seedLength: cut,
  1854. ...forkComposition.agentPreset === undefined
  1855. ? {}
  1856. : { agentPreset: forkComposition.agentPreset },
  1857. },
  1858. agentOptions,
  1859. setup: forkComposition.setup,
  1860. })
  1861. } catch (error: unknown) {
  1862. return err(request, {
  1863. code: 'internal',
  1864. message: `failed to fork session "${sessionId}": ${String(error)}`,
  1865. details: {},
  1866. })
  1867. }
  1868. // An ordinary source keeps its direct Workspace. A subagent source is
  1869. // not listed there, so its ordinary fork joins the nearest owning
  1870. // ancestor instead. The child is already published if attach fails.
  1871. if (workspace !== undefined) {
  1872. try {
  1873. await workspace.attachSession(childId)
  1874. } catch (error: unknown) {
  1875. return err(request, {
  1876. code: 'workspace-attach-failed',
  1877. message: `session "${childId}" was forked but could not attach to workspace "${workspace.id}": ${String(error)}`,
  1878. details: { sessionId: childId, workspaceId: workspace.id },
  1879. })
  1880. }
  1881. }
  1882. return ok(request, { sessionId: childId })
  1883. },
  1884. async prompt(request) {
  1885. const { sessionId, mode, content } = request.payload
  1886. const found = await agentFor(sessionId)
  1887. if ('error' in found) return err(request, found.error)
  1888. const agent = found.agent
  1889. // The rpcId rides MessageSource into user/message (merge declaration in api/sessions.ts; provisional correlation).
  1890. const source: MessageSource = { kind: 'user', rpcId: request.rpcId }
  1891. try {
  1892. const message: UserMessage = createUserMessage({ content, source })
  1893. if (mode === 'steer') agent.steer(message)
  1894. else agent.followup(message)
  1895. } catch (error: unknown) {
  1896. // A synchronous throw from steer/followup means disposed or invalid input; surface as agent-busy with the reason attached.
  1897. return err(request, { code: 'agent-busy', message: 'prompt rejected', details: { reason: String(error) } })
  1898. }
  1899. return ok(request, { accepted: true as const })
  1900. },
  1901. updateQueue(request) {
  1902. const { sessionId, itemId, action } = request.payload
  1903. const agent = ctx.agents.get(sessionId)
  1904. if (agent !== undefined && hasSubagentOwner(agent.session, agent)) {
  1905. return Promise.resolve(err(request, subagentOwnershipError(sessionId)))
  1906. }
  1907. if (agent === undefined) {
  1908. return Promise.resolve(err(request, {
  1909. code: 'queue-item-not-found',
  1910. message: 'queued item is no longer pending',
  1911. details: { itemId },
  1912. }))
  1913. }
  1914. const target = agent.inbox.nextTurn.some(message => message.id === itemId)
  1915. ? 'next-turn'
  1916. : agent.inbox.nextStep.some(message => message.id === itemId) ? 'next-step' : undefined
  1917. const message = target === undefined
  1918. ? undefined
  1919. : (target === 'next-turn' ? agent.inbox.nextTurn : agent.inbox.nextStep)
  1920. .find(candidate => candidate.id === itemId)
  1921. if (target === undefined || message === undefined) {
  1922. return Promise.resolve(err(request, {
  1923. code: 'queue-item-not-found',
  1924. message: 'queued item is no longer pending',
  1925. details: { itemId },
  1926. }))
  1927. }
  1928. if (action.kind === 'steer' && (target !== 'next-turn' || agent.status !== 'running')) {
  1929. return Promise.resolve(err(request, {
  1930. code: 'steer-unavailable',
  1931. message: 'current turn no longer accepts steering',
  1932. details: { itemId },
  1933. }))
  1934. }
  1935. if (action.kind === 'edit') {
  1936. agent.inbox.replace(itemId, freezeMessage({ ...message, content: action.content }))
  1937. } else {
  1938. agent.inbox.remove(itemId)
  1939. if (action.kind === 'steer') agent.steer(message)
  1940. }
  1941. return Promise.resolve(ok(request, { accepted: true as const }))
  1942. },
  1943. cancel(request) {
  1944. const { sessionId } = request.payload
  1945. const agent = ctx.agents.get(sessionId)
  1946. if (agent === undefined) {
  1947. return Promise.resolve(err(request, {
  1948. code: 'session-not-found',
  1949. message: `session "${sessionId}" not found (not attached)`,
  1950. details: { sessionId },
  1951. }))
  1952. }
  1953. if (hasSubagentOwner(agent.session, agent)) {
  1954. return Promise.resolve(err(request, subagentOwnershipError(sessionId)))
  1955. }
  1956. agent.cancel({ kind: 'user' }, { keepInbox: true })
  1957. return Promise.resolve(ok(request, { accepted: true as const }))
  1958. },
  1959. },
  1960. subagents: {
  1961. async list(request, signal) {
  1962. try {
  1963. const entries = await ctx.subagents.listChildren(request.payload.parentSessionId, signal)
  1964. return ok(request, {
  1965. entries: entries.map(entry => entry.kind === 'child'
  1966. ? {
  1967. ...entry,
  1968. activity: ctx.agents.get(entry.id)?.status === 'running' ? 'running' : 'inactive',
  1969. }
  1970. : entry),
  1971. parentAvailable: ctx.agents.get(request.payload.parentSessionId) !== undefined,
  1972. })
  1973. } catch (error: unknown) {
  1974. if (signal?.aborted
  1975. || (error instanceof SubagentError && error.code === 'CANCELLED')
  1976. || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')) {
  1977. return err(request, {
  1978. code: 'cancelled',
  1979. message: 'subagent catalog read was cancelled',
  1980. details: {},
  1981. })
  1982. }
  1983. return err(request, {
  1984. code: 'internal',
  1985. message: 'subagent catalog read failed',
  1986. details: {},
  1987. })
  1988. }
  1989. },
  1990. async history(request, signal) {
  1991. const {
  1992. parentSessionId, childSessionId, mode, beforeSeq, maxMessages,
  1993. } = request.payload
  1994. const verified = await catalogChild(ctx, {
  1995. parentSessionId, childSessionId, mode,
  1996. }, signal)
  1997. if (verified.error !== undefined) return err(request, verified.error)
  1998. try {
  1999. const snapshot = await ctx.sessionQuery.readSession(childSessionId)
  2000. signal?.throwIfAborted()
  2001. if (snapshot.session.parentSession !== parentSessionId) {
  2002. return err(request, {
  2003. code: 'subagent-unauthorized',
  2004. message: 'subagent parent changed during history read',
  2005. details: { childSessionId },
  2006. })
  2007. }
  2008. const page = historyPage(ctx, snapshot.events, beforeSeq, maxMessages, ctx.agents.get(childSessionId))
  2009. const projections = beforeSeq === undefined
  2010. ? detachedProjectionsFor(ctx, snapshot.events)
  2011. : undefined
  2012. return ok(request, { ...page, ...projections === undefined ? {} : { projections } })
  2013. } catch (error: unknown) {
  2014. if (signal?.aborted
  2015. || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')) {
  2016. return err(request, {
  2017. code: 'cancelled',
  2018. message: 'subagent history read was cancelled',
  2019. details: {},
  2020. })
  2021. }
  2022. if (error instanceof SessionQueryError
  2023. && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') {
  2024. return err(request, {
  2025. code: 'subagent-not-found',
  2026. message: 'subagent disappeared during history read',
  2027. details: { parentSessionId, childSessionId },
  2028. })
  2029. }
  2030. return err(request, {
  2031. code: 'internal',
  2032. message: 'subagent history read failed',
  2033. details: {},
  2034. })
  2035. }
  2036. },
  2037. async prompt(request, signal) {
  2038. const { parentSessionId, childSessionId, content } = request.payload
  2039. const parent = ctx.agents.get(parentSessionId)
  2040. if (parent === undefined) {
  2041. return err(request, {
  2042. code: 'subagent-parent-unavailable',
  2043. message: `parent session "${parentSessionId}" is not live`,
  2044. details: { parentSessionId },
  2045. })
  2046. }
  2047. const verified = await catalogChild(ctx, {
  2048. parentSessionId, childSessionId, mode: 'continuable',
  2049. }, signal)
  2050. if (verified.error !== undefined) return err(request, verified.error)
  2051. try {
  2052. const messageId = await ctx.subagents.followup(parent, childSessionId, content, {
  2053. source: { kind: 'user', rpcId: request.rpcId },
  2054. signal,
  2055. })
  2056. return ok(request, { messageId })
  2057. } catch (error: unknown) {
  2058. return subagentPromptError(request, error, signal)
  2059. }
  2060. },
  2061. },
  2062. workspace: {
  2063. list(request) {
  2064. return Promise.resolve(ok(request, {
  2065. items: ctx.workspace.list().map(workspaceView),
  2066. archivedSessionIds: [...ctx.workspace.archivedSessionIds],
  2067. }))
  2068. },
  2069. // Exactly one of path/name arrives (schema refine). Existing-folder
  2070. // adoption reuses its canonical path; create-by-name rejects a name
  2071. // already present in the registry.
  2072. // TODO: the create-by-name branch lost its last product consumer when
  2073. // the Web picker collapsed onto the directory flow
  2074. // (.agents/notes/implemented/simplification/2026-07-31-one-route-to-add-a-workspace.md).
  2075. // Delete it with the wire schema's `name` member, this
  2076. // `defaults.workspaceRoot`, the client seam that carried the name
  2077. // (`WorkspaceCreateInput`, `WorkspacesService.create`'s `{ name }` arm,
  2078. // `intentName`'s name branch, the manager's "name under workspaceRoot"
  2079. // contract), and the `dsh web --workspace-root` flag plus its apps/cli
  2080. // README lines, which exist only to feed it.
  2081. async create(request) {
  2082. const { payload } = request
  2083. let path: string
  2084. if (payload.name !== undefined) {
  2085. const name = payload.name.trim()
  2086. if (name === '' || name === '.' || name === '..' || /[/\\]/.test(name)) {
  2087. return err(request, {
  2088. code: 'workspace-invalid-path',
  2089. message: `workspace name must be one non-empty path segment, got "${payload.name}"`,
  2090. details: { path: payload.name },
  2091. })
  2092. }
  2093. path = join(defaults.workspaceRoot, name)
  2094. } else {
  2095. path = payload.path as string
  2096. }
  2097. try {
  2098. const name = payload.name?.trim()
  2099. const { workspace, created } = await ensureWorkspace(
  2100. path,
  2101. name,
  2102. name !== undefined,
  2103. name !== undefined,
  2104. )
  2105. return ok(request, { workspace: workspaceView(workspace), created })
  2106. } catch (error: unknown) {
  2107. if (error instanceof WorkspaceNameConflictError) {
  2108. return err(request, {
  2109. code: 'workspace-name-conflict',
  2110. message: error.message,
  2111. details: { name: error.workspaceName },
  2112. })
  2113. }
  2114. if (error instanceof WorkspaceDirectoryCreationError) {
  2115. return err(request, { code: 'internal', message: error.message, details: {} })
  2116. }
  2117. // The registry rejects a path that does not resolve to an existing
  2118. // directory (realpath ENOENT / not-a-directory) — the business
  2119. // error of the typed-path flow, surfaced as a validation failure.
  2120. return err(request, {
  2121. code: 'workspace-invalid-path',
  2122. message: `cannot create a workspace at "${path}": ${error instanceof Error ? error.message : String(error)}`,
  2123. details: { path },
  2124. })
  2125. }
  2126. },
  2127. async rename(request) {
  2128. const { payload } = request
  2129. const workspace = ctx.workspace.get(brandWorkspaceId(payload.workspaceId))
  2130. if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
  2131. const title = payload.title.trim()
  2132. // Uniqueness AND the same-title no-op both ride the create chain so
  2133. // they observe the state left by earlier queued renames — checked
  2134. // up front, a queued A→A could report success while an earlier A→B
  2135. // still lands afterwards.
  2136. const operation = workspaceCreationChain.then(async () => {
  2137. if (title === workspace.title) return
  2138. if (ctx.workspace.list().some(other => other.id !== workspace.id && other.title === title)) {
  2139. throw new WorkspaceNameConflictError(title)
  2140. }
  2141. await workspace.setTitle(title)
  2142. })
  2143. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  2144. try {
  2145. await operation
  2146. } catch (error: unknown) {
  2147. if (error instanceof WorkspaceNameConflictError) {
  2148. return err(request, {
  2149. code: 'workspace-name-conflict',
  2150. message: error.message,
  2151. details: { name: error.workspaceName },
  2152. })
  2153. }
  2154. throw error
  2155. }
  2156. return ok(request, { workspace: workspaceView(workspace) })
  2157. },
  2158. async delete(request) {
  2159. const { workspaceId } = request.payload
  2160. const operation = workspaceCreationChain.then(() =>
  2161. ctx.workspace.delete(brandWorkspaceId(workspaceId)))
  2162. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  2163. if (!await operation) return workspaceNotFound(request, workspaceId)
  2164. return ok(request, { deleted: true as const })
  2165. },
  2166. async insertSessionBefore(request) {
  2167. const { payload } = request
  2168. const workspace = ctx.workspace.get(brandWorkspaceId(payload.workspaceId))
  2169. if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
  2170. try {
  2171. await workspace.insertSessionBefore(payload.sessionId, payload.beforeSessionId)
  2172. } catch (error: unknown) {
  2173. // Only the entity's unaccounted-id rejection is the business code;
  2174. // storage/durability failures propagate as internal errors.
  2175. if (!(error instanceof WorkspaceMoveInvalidError)) throw error
  2176. return err(request, {
  2177. code: 'workspace-move-invalid',
  2178. message: error.message,
  2179. details: {
  2180. workspaceId: payload.workspaceId,
  2181. sessionId: payload.sessionId,
  2182. ...payload.beforeSessionId === undefined ? {} : { beforeSessionId: payload.beforeSessionId },
  2183. },
  2184. })
  2185. }
  2186. return ok(request, { workspace: workspaceView(workspace) })
  2187. },
  2188. async archiveSession(request) {
  2189. const { sessionId } = request.payload
  2190. try {
  2191. await ctx.workspace.archiveSession(sessionId)
  2192. } catch (error: unknown) {
  2193. // Only the registry's unknown-session rejection is the business
  2194. // code; storage/durability failures propagate as internal errors.
  2195. if (!(error instanceof WorkspaceUnknownSessionError)) throw error
  2196. return err(request, {
  2197. code: 'session-not-found',
  2198. message: error.message,
  2199. details: { sessionId },
  2200. })
  2201. }
  2202. return ok(request, { archivedSessionIds: [...ctx.workspace.archivedSessionIds] })
  2203. },
  2204. },
  2205. host: {
  2206. describe(request) {
  2207. // TODO(step2): version should read apps/cli's package.json; placeholder for now.
  2208. return Promise.resolve(ok(request, {
  2209. version: '0.0.1',
  2210. // Same source as session.create's fallback: the UI's default project
  2211. // must match where an unspecified-cwd session actually lands.
  2212. cwd: defaults.cwd,
  2213. provider: defaults.provider,
  2214. model: defaults.model,
  2215. attachedSessions: ctx.agents.list().length,
  2216. }))
  2217. },
  2218. async pickDirectory(request, signal) {
  2219. const capability = ctx.directoryPicker.capability()
  2220. if (capability.kind !== 'native') {
  2221. return err(request, {
  2222. code: 'directory-picker-unavailable',
  2223. message: `host.pickDirectory needs the native capability; the composed picker serves "${capability.kind}"`,
  2224. details: { capability: capability.kind },
  2225. })
  2226. }
  2227. try {
  2228. const path = await capability.pick(signal)
  2229. return ok(request, { path })
  2230. } catch (error: unknown) {
  2231. if (signal.aborted) {
  2232. return err(request, {
  2233. code: 'cancelled',
  2234. message: 'directory picker was aborted',
  2235. details: {},
  2236. })
  2237. }
  2238. return err(request, {
  2239. code: 'internal',
  2240. message: `directory picker failed: ${error instanceof Error ? error.message : String(error)}`,
  2241. details: {},
  2242. })
  2243. }
  2244. },
  2245. async listDirectory(request, signal) {
  2246. const capability = ctx.directoryPicker.capability()
  2247. if (capability.kind !== 'browse') {
  2248. return err(request, {
  2249. code: 'directory-picker-unavailable',
  2250. message: `host.listDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
  2251. details: { capability: capability.kind },
  2252. })
  2253. }
  2254. try {
  2255. // The carrier's signal follows the caller: a disconnect or timeout
  2256. // stops the backend's directory scan instead of outliving it.
  2257. return ok(request, await capability.list(request.payload.path, signal))
  2258. } catch (error: unknown) {
  2259. // An abort is the caller's own timeout/disconnect, not a server
  2260. // failure — same code pickDirectory and command.execute report.
  2261. if (signal.aborted) {
  2262. return err(request, { code: 'cancelled', message: 'directory listing was aborted', details: {} })
  2263. }
  2264. return err(request, directoryError(error))
  2265. }
  2266. },
  2267. async createDirectory(request) {
  2268. const capability = ctx.directoryPicker.capability()
  2269. if (capability.kind !== 'browse') {
  2270. return err(request, {
  2271. code: 'directory-picker-unavailable',
  2272. message: `host.createDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
  2273. details: { capability: capability.kind },
  2274. })
  2275. }
  2276. try {
  2277. return ok(request, { path: await capability.createDirectory(request.payload.path, request.payload.name) })
  2278. } catch (error: unknown) {
  2279. return err(request, directoryError(error))
  2280. }
  2281. },
  2282. async openPath(request, signal) {
  2283. return openPath(request, request.payload.path, signal)
  2284. },
  2285. },
  2286. commands: {
  2287. // Both methods address one session's agent. agentFor resumes on miss
  2288. // and fences every subagent-owned identity with `agent-busy`; the
  2289. // api/commands.ts module contract owns that fence's wording, so this
  2290. // comment only notes the routing shape: clients send a sessionId for a
  2291. // published session, and resume restores an existing entity.
  2292. async list(request) {
  2293. // Missing service = the deployment omitted dsh-commands from its
  2294. // composition, not an empty catalog: fail loud instead of serving [].
  2295. const commands = ctx.get('commands')
  2296. if (commands === undefined) {
  2297. 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: {} })
  2298. }
  2299. const found = await agentFor(request.payload.sessionId)
  2300. if ('error' in found) return err(request, found.error)
  2301. return ok(request, { commands: commands.list(found.agent) })
  2302. },
  2303. async execute(request, signal) {
  2304. const commands = ctx.get('commands')
  2305. if (commands === undefined) {
  2306. 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: {} })
  2307. }
  2308. const { sessionId, line } = request.payload
  2309. const found = await agentFor(sessionId)
  2310. if ('error' in found) return err(request, found.error)
  2311. try {
  2312. // Pure admission: the executor's durable command/run + command/done
  2313. // pair (broadcast on the mux stream) carries the outcome; the
  2314. // response reports whether the line resolved to a handler, plus the
  2315. // minted pairing id so the issuing client can correlate its request
  2316. // with the flow node the lifecycle events produce.
  2317. const execution = await commands.execute(found.agent, line, signal)
  2318. return ok(request, execution === undefined
  2319. ? { matched: false }
  2320. : { matched: true, commandId: execution.commandId })
  2321. } catch (error: unknown) {
  2322. if (signal.aborted) return err(request, { code: 'cancelled', message: 'command execution was aborted', details: {} })
  2323. return err(request, { code: 'internal', message: `command failed: ${String(error)}`, details: {} })
  2324. }
  2325. },
  2326. },
  2327. goals: {
  2328. // Mutations only — the read side is the 'goal' session projection.
  2329. // Every verb resolves the session's agent (agentFor: implicit cold
  2330. // resume, the command.* precedent) and acknowledges with the new CAS
  2331. // ref; the committed goal/change event carries the whole value to every
  2332. // client through the projection frames.
  2333. async create(request) {
  2334. const { objective, maxGoalRounds } = request.payload
  2335. return mutateGoal(request, (goals, agent) => goals.create(agent, {
  2336. objective,
  2337. ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
  2338. }))
  2339. },
  2340. async edit(request) {
  2341. const { ref, objective, maxGoalRounds } = request.payload
  2342. return mutateGoal(request, (goals, agent) => goals.edit(agent, ref, {
  2343. ...(objective !== undefined ? { objective } : {}),
  2344. ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
  2345. }))
  2346. },
  2347. async pause(request) {
  2348. return mutateGoal(request, (goals, agent) => goals.pause(agent, request.payload.ref))
  2349. },
  2350. async resume(request) {
  2351. return mutateGoal(request, (goals, agent) => goals.resume(agent, request.payload.ref))
  2352. },
  2353. async complete(request) {
  2354. return mutateGoal(request, (goals, agent) => goals.complete(agent, request.payload.ref))
  2355. },
  2356. async clear(request) {
  2357. const found = await agentFor(request.payload.sessionId)
  2358. if ('error' in found) return err(request, found.error)
  2359. const goals = goalServiceFor(found.agent)
  2360. if ('error' in goals) return err(request, goals.error)
  2361. try {
  2362. goals.clear(found.agent, request.payload.ref)
  2363. return ok(request, { cleared: true as const })
  2364. } catch (error: unknown) {
  2365. return goalError(request, error)
  2366. }
  2367. },
  2368. },
  2369. agentPresets: {
  2370. // A deployment with no roster answers with an empty list rather than an
  2371. // error: composing no presets is a valid deployment, and the browser
  2372. // simply offers no choice.
  2373. async list(request) {
  2374. const presets = ctx.get('agentPresets')
  2375. if (presets === undefined) return ok(request, { presets: [] })
  2376. const defaultId = presets.defaultId
  2377. return ok(request, {
  2378. presets: (await presets.list()).map(preset => ({
  2379. id: preset.id,
  2380. trust: preset.trust,
  2381. isDefault: preset.id === defaultId,
  2382. })),
  2383. })
  2384. },
  2385. // Recomposing is limited to a blank session because a started
  2386. // conversation's history was produced under its preset's tools; the
  2387. // agent and the session survive, only the composition is swapped.
  2388. async select(request) {
  2389. const { sessionId, agentPreset } = request.payload
  2390. const presets = ctx.get('agentPresets')
  2391. if (presets === undefined) {
  2392. return err(request, {
  2393. code: 'agent-preset-not-found',
  2394. message: 'this deployment composes no agent presets',
  2395. details: { agentPreset, available: [] },
  2396. })
  2397. }
  2398. const found = await agentFor(sessionId)
  2399. if ('error' in found) return err(request, found.error)
  2400. const { agent } = found
  2401. if (!sessionBlank(agent.session)) {
  2402. return err(request, {
  2403. code: 'agent-preset-locked',
  2404. message: `session "${sessionId}" has already started; its agent preset is fixed`,
  2405. details: { sessionId, agentPreset },
  2406. })
  2407. }
  2408. try {
  2409. const preset = await presets.recompose(agent.ctx, agentPreset)
  2410. return ok(request, { agentPreset: preset.id })
  2411. } catch (error: unknown) {
  2412. if (error instanceof UnknownPresetError) {
  2413. return err(request, {
  2414. code: 'agent-preset-not-found',
  2415. message: error.message,
  2416. details: { agentPreset: error.presetId, available: [...error.available] },
  2417. })
  2418. }
  2419. if (error instanceof PresetMountError) {
  2420. return err(request, {
  2421. code: 'agent-preset-invalid',
  2422. message: error.message,
  2423. details: { agentPreset: error.presetId, reason: error.reason },
  2424. })
  2425. }
  2426. return err(request, {
  2427. code: 'internal',
  2428. message: `failed to select agent preset "${agentPreset}": ${String(error)}`,
  2429. details: {},
  2430. })
  2431. }
  2432. },
  2433. },
  2434. skills: {
  2435. // Skill lookup never touches the Agent registry: the session address
  2436. // resolves to a canonical cwd from the host-resident session header, so
  2437. // listing skills cannot create or resume an agent as a side effect.
  2438. async list(request) {
  2439. const { sessionId } = request.payload
  2440. const session = ctx.sessions.get(sessionId)
  2441. if (session === undefined) {
  2442. return err(request, {
  2443. code: 'session-not-found',
  2444. message: `session "${sessionId}" not found (not attached)`,
  2445. details: { sessionId },
  2446. })
  2447. }
  2448. if (session.header.cwd === undefined) {
  2449. // Every served session records its project at create time; a
  2450. // cwd-less header is a pre-project legacy log (not served).
  2451. return err(request, { code: 'internal', message: `session "${sessionId}" has no project cwd`, details: {} })
  2452. }
  2453. const cwd = session.header.cwd
  2454. // The registry is per session when a preset mounts one — a preset
  2455. // ships its own skill directory, so the catalog IS the session's — and
  2456. // that instance sits behind an `isolate` realm no host context
  2457. // resolves. Address it through the live agent; `agents.get` keeps the
  2458. // no-side-effect stance above (a cold session creates nothing and
  2459. // falls through to whatever the host composes).
  2460. const live = ctx.agents.get(sessionId)
  2461. const presets = ctx.get('agentPresets')
  2462. const scoped = live === undefined ? undefined : presets?.serviceFor(live, 'skills')
  2463. // Same stance as the commands domain: a missing service means no
  2464. // composition mounts dsh-skill, not an empty catalog. `ctx.get` also
  2465. // keeps this handler independent of the gateway plugin's inject list
  2466. // (an undeclared `ctx.skills` property read fails the reflect proxy).
  2467. const skillRegistry = scoped ?? ctx.get('skills')
  2468. if (skillRegistry === undefined) {
  2469. return err(request, { code: 'internal', message: 'skill registry is absent: neither this session\'s agent preset nor the host composition mounts @deepseek-ai/dsh-skill', details: {} })
  2470. }
  2471. try {
  2472. const skills = (await skillRegistry.list({ cwd }))
  2473. .filter(skill => skill.invocation.modelInvocable && skill.invocation.userInvocable)
  2474. return ok(request, {
  2475. skills: skills.map(skill => ({
  2476. name: skill.name,
  2477. description: skill.description,
  2478. ...skill.whenToUse === undefined ? {} : { whenToUse: skill.whenToUse },
  2479. })),
  2480. })
  2481. } catch (error: unknown) {
  2482. return err(request, { code: 'internal', message: `skill listing failed: ${String(error)}`, details: {} })
  2483. }
  2484. },
  2485. },
  2486. settings: {
  2487. describe(request) {
  2488. const settings = ctx.get('settings')
  2489. if (settings === undefined) return Promise.resolve(err(request, settingsAbsent()))
  2490. const exposed = exposedNamespaces()
  2491. return Promise.resolve(ok(request, {
  2492. writable: settings.writable,
  2493. hasDocument: settings.documentPath !== undefined,
  2494. namespaces: settings.describe({ redactSecrets: true })
  2495. .filter(descriptor => exposed.has(String(descriptor.ns)))
  2496. .map(namespaceView),
  2497. }))
  2498. },
  2499. async openDocument(request, signal) {
  2500. const settings = ctx.get('settings')
  2501. if (settings === undefined) return err(request, settingsAbsent())
  2502. if (isAborted(signal)) {
  2503. return err(request, {
  2504. code: 'cancelled',
  2505. message: 'settings document open was aborted',
  2506. details: {},
  2507. })
  2508. }
  2509. let path: string | undefined
  2510. try {
  2511. path = await settings.prepareDocument()
  2512. } catch (error: unknown) {
  2513. if (isAborted(signal)) {
  2514. return err(request, {
  2515. code: 'cancelled',
  2516. message: 'settings document preparation was aborted',
  2517. details: {},
  2518. })
  2519. }
  2520. return err(request, {
  2521. code: 'internal',
  2522. message: `settings document preparation failed: ${error instanceof Error ? error.message : String(error)}`,
  2523. details: {},
  2524. })
  2525. }
  2526. if (path === undefined) {
  2527. return err(request, {
  2528. code: 'internal',
  2529. message: 'settings provider has no local document to open',
  2530. details: {},
  2531. })
  2532. }
  2533. if (isAborted(signal)) {
  2534. return err(request, {
  2535. code: 'cancelled',
  2536. message: 'settings document open was aborted',
  2537. details: {},
  2538. })
  2539. }
  2540. return openTextFile(request, path, signal)
  2541. },
  2542. update: request => settingsWrite(request, request.payload.ns, 'update', request.payload.patch, request.payload.expectedRevision),
  2543. replace: request => settingsWrite(request, request.payload.ns, 'replace', request.payload.section, request.payload.expectedRevision),
  2544. mutate: request => settingsWrite(request, request.payload.ns, 'mutate', request.payload.ops, request.payload.expectedRevision),
  2545. },
  2546. credentials: {
  2547. async describe(request) {
  2548. const credentials = ctx.get('credentials')
  2549. if (credentials === undefined) return err(request, credentialsAbsent())
  2550. const entries = await Promise.all(request.payload.refs.map(async (ref) => {
  2551. const info = await credentials.describe(credentialRef(ref))
  2552. const view: CredentialView = {
  2553. configured: info.configured,
  2554. ...info.source === undefined ? {} : { source: info.source },
  2555. writable: info.writable,
  2556. }
  2557. return [ref, view] as const
  2558. }))
  2559. return ok(request, { credentials: Object.fromEntries(entries) })
  2560. },
  2561. async set(request) {
  2562. const credentials = ctx.get('credentials')
  2563. if (credentials === undefined) return err(request, credentialsAbsent())
  2564. const { ref, value } = request.payload
  2565. try {
  2566. await credentials.set(credentialRef(ref), value)
  2567. } catch (error: unknown) {
  2568. return err(request, {
  2569. code: 'credential-rejected',
  2570. message: error instanceof Error ? error.message : String(error),
  2571. details: { ref },
  2572. })
  2573. }
  2574. return ok(request, {})
  2575. },
  2576. async unset(request) {
  2577. const credentials = ctx.get('credentials')
  2578. if (credentials === undefined) return err(request, credentialsAbsent())
  2579. const { ref } = request.payload
  2580. try {
  2581. await credentials.unset(credentialRef(ref))
  2582. } catch (error: unknown) {
  2583. return err(request, {
  2584. code: 'credential-rejected',
  2585. message: error instanceof Error ? error.message : String(error),
  2586. details: { ref },
  2587. })
  2588. }
  2589. return ok(request, {})
  2590. },
  2591. },
  2592. llm: {
  2593. providers(request) {
  2594. const registered = ctx.llm.listProviders()
  2595. const active = new Set(registered.map(provider => provider.id))
  2596. const directory = ctx.llm.listConfigurableProviders()
  2597. const declared = new Set(directory.map(entry => entry.provider))
  2598. const views = directory.map(entry => ({
  2599. provider: entry.provider,
  2600. displayName: entry.displayName,
  2601. settingsNs: entry.settingsNs,
  2602. settingsPath: [...entry.settingsPath],
  2603. active: active.has(entry.provider),
  2604. }))
  2605. // Routes registered without a directory declaration still appear —
  2606. // they exist and serve models — just with no settings address.
  2607. for (const provider of registered) {
  2608. if (declared.has(provider.id)) continue
  2609. views.push({
  2610. provider: provider.id,
  2611. displayName: provider.name,
  2612. settingsNs: '',
  2613. settingsPath: [],
  2614. active: true,
  2615. })
  2616. }
  2617. return Promise.resolve(ok(request, { providers: views }))
  2618. },
  2619. async models(request) {
  2620. return ok(request, await buildModelCatalog(ctx))
  2621. },
  2622. async discoverModels(request, signal) {
  2623. const { settingsNs, provider, baseURL, api, apiKey } = request.payload
  2624. try {
  2625. const models = await ctx.llm.discoverModels(settingsNs, {
  2626. ...provider === undefined ? {} : { provider },
  2627. ...baseURL === undefined ? {} : { baseURL },
  2628. ...api === undefined ? {} : { api },
  2629. ...apiKey === undefined ? {} : { apiKey },
  2630. ...signal === undefined ? {} : { signal },
  2631. })
  2632. return ok(request, { models })
  2633. } catch (error: unknown) {
  2634. // Every failure here is the user's next move, not a transport fault:
  2635. // a wrong endpoint, a rejected key, or a protocol with no listing all
  2636. // end at the same place — fill the models in by hand. The details
  2637. // repeat only what the caller already sent, never the credential.
  2638. return err(request, {
  2639. code: 'model-discovery-failed',
  2640. message: error instanceof Error ? error.message : String(error),
  2641. details: { settingsNs, ...baseURL === undefined ? {} : { baseURL } },
  2642. })
  2643. }
  2644. },
  2645. },
  2646. events: {
  2647. mux(_request, signal) {
  2648. const queue = new FrameQueue<RpcRequest<MuxFrame>>()
  2649. muxQueues.add(queue)
  2650. for (const session of ctx.sessions.list()) {
  2651. subscribeSession(queue, session)
  2652. }
  2653. for (const pending of pendingQuestions.values()) {
  2654. queue.push({
  2655. rpcId: pending.rpcId,
  2656. payload: {
  2657. type: 'question/requested', sessionId: pending.sessionId,
  2658. questions: pending.questions,
  2659. },
  2660. })
  2661. }
  2662. // Refresh recovery: still-pending approval questions replay with their
  2663. // stable rpcId so a reconnecting client can still answer them.
  2664. for (const pending of pendingApprovals.values()) queue.push(requestedFrame(pending))
  2665. // Queue snapshot baseline (pendingQuestions precedent): frames replayed
  2666. // in arrival order per session; a reconnecting client rebuilds its
  2667. // queue view from these alone.
  2668. for (const session of ctx.sessions.list()) {
  2669. const agent = ctx.agents.get(session.id)
  2670. if (agent?.session === session && agent.inbox.hasPending) {
  2671. queue.push(frame({ type: 'session/queue', sessionId: session.id, items: queueItems(agent) }))
  2672. }
  2673. }
  2674. // Per-session open-call table for result-view pairing. Bounded by the
  2675. // per-turn call count: entries clear on turn/end; a table miss (stream
  2676. // opened mid-turn) backscans the session's in-memory events instead.
  2677. const openCalls = new Map<SessionId, Map<string, { name: string; args: unknown }>>()
  2678. const disposers = [
  2679. ctx.on('session/event', (session: Session, event: SessionEvent) => {
  2680. if (event.type === 'tool/call') {
  2681. const data = event.data as ToolCallData
  2682. try {
  2683. let table = openCalls.get(session.id)
  2684. if (table === undefined) openCalls.set(session.id, table = new Map<string, { name: string; args: unknown }>())
  2685. table.set(data.callId, { name: data.name, args: JSON.parse(data.arguments) })
  2686. } catch {
  2687. // Unparseable model arguments: leave the table unset; the result view soft-falls.
  2688. }
  2689. } else if (event.type === 'turn/end') {
  2690. openCalls.delete(session.id)
  2691. }
  2692. const view = viewFor(
  2693. ctx, event,
  2694. callId => openCalls.get(session.id)?.get(callId) ?? backscanArgs(session.events, callId),
  2695. ctx.agents.get(session.id),
  2696. )
  2697. queue.push(frame({ type: 'session/event', sessionId: session.id, event, ...view === undefined ? {} : { view } }))
  2698. }),
  2699. ctx.on('session/created', (session: Session) => {
  2700. subscribeSession(queue, session)
  2701. }),
  2702. ctx.on('session/disposed', (session: Session) => {
  2703. openCalls.delete(session.id)
  2704. }),
  2705. ]
  2706. return queue.iterate(signal, () => {
  2707. muxQueues.delete(queue)
  2708. for (const dispose of disposers) dispose()
  2709. })
  2710. },
  2711. host(_request, signal) {
  2712. const queue = new FrameQueue<RpcRequest<HostFrame>>()
  2713. const committedWorkspaceIds = new Set(
  2714. ctx.workspace.list().map(workspace => String(workspace.id)),
  2715. )
  2716. // Frame-dedup baseline, same posture as committedWorkspaceIds: the
  2717. // stream opens against the current set; workspace.list re-baselines
  2718. // reconnecting clients, so only later changes need frames.
  2719. let archivedSessionIds = ctx.workspace.archivedSessionIds
  2720. const disposers = [
  2721. ctx.on('session/created', (session: Session) => {
  2722. queue.push(frame({
  2723. type: 'host/session-added',
  2724. sessionId: session.id,
  2725. // Derived at frame time like summarize(); a just-created session
  2726. // has run no turn yet, so this is constantly true in practice.
  2727. blank: sessionBlank(session),
  2728. // Including cwd lets the client group the new session without refreshing the list.
  2729. ...sessionListFields(session.header),
  2730. }))
  2731. }),
  2732. ctx.on('session/disposed', (session: Session) => {
  2733. queue.push(frame({ type: 'host/session-removed', sessionId: session.id }))
  2734. }),
  2735. ctx.on('agent/status', ({ agent, status }: { agent: Agent; status: AgentStatus }) => {
  2736. queue.push(frame({ type: 'host/session-status', sessionId: agent.id, running: status === 'running' }))
  2737. }),
  2738. ctx.on('agent/error', ({ agent, error }: { agent: Agent; error: unknown }) => {
  2739. queue.push(frame({ type: 'host/agent-error', sessionId: agent.id, message: errorChain(error) }))
  2740. }),
  2741. ctx.on('domain/changed', (change) => {
  2742. if (change.domain !== 'workspace') return
  2743. if (change.table === '') {
  2744. if (change.operation !== 'put') return
  2745. const state = workspaceDomainState.parse(change.value)
  2746. for (const workspaceId of state.workspaceIds) {
  2747. if (committedWorkspaceIds.has(workspaceId)) continue
  2748. const workspace = ctx.workspace.get(workspaceId)
  2749. if (workspace === undefined) {
  2750. throw new Error(`committed workspace registry references missing workspace "${workspaceId}"`)
  2751. }
  2752. committedWorkspaceIds.add(workspaceId)
  2753. queue.push(frame({ type: 'host/workspace-changed', workspace: workspaceView(workspace) }))
  2754. }
  2755. if (state.archivedSessionIds.length !== archivedSessionIds.length
  2756. || state.archivedSessionIds.some((id, index) => id !== archivedSessionIds[index])) {
  2757. archivedSessionIds = state.archivedSessionIds
  2758. queue.push(frame({
  2759. type: 'host/archived-sessions-changed',
  2760. archivedSessionIds: [...state.archivedSessionIds],
  2761. }))
  2762. }
  2763. return
  2764. }
  2765. if (change.table !== 'workspaces') return
  2766. if (change.operation === 'deleted') {
  2767. if (!committedWorkspaceIds.delete(change.key)) return
  2768. queue.push(frame({
  2769. type: 'host/workspace-removed',
  2770. workspaceId: change.key as WorkspaceId,
  2771. }))
  2772. return
  2773. }
  2774. if (!committedWorkspaceIds.has(change.key)) return
  2775. // Existing-entity table writes are complete attach/touch commits.
  2776. // A new entity's first put waits for the global registry write above.
  2777. queue.push(frame({
  2778. type: 'host/workspace-changed',
  2779. workspace: changedWorkspaceView(change.key, change.value),
  2780. }))
  2781. }),
  2782. ctx.on('commands/change', () => {
  2783. queue.push(frame({ type: 'host/commands-changed' }))
  2784. }),
  2785. ctx.on('settings/document-updated', (ns) => {
  2786. // The RAW-section event, not the resolved one: a field going from
  2787. // inherited to overridden leaves the resolved value equal, and a
  2788. // configuration client still has to re-read (its held revision is
  2789. // stale, and the field's meaning changed).
  2790. const name = String(ns)
  2791. queue.push(frame({ type: 'host/settings-changed', ns: name }))
  2792. // A provider's own settings carry its model catalog and endpoint,
  2793. // so a change there invalidates the model list even when the route
  2794. // set is untouched — `llm/adapters-updated` alone misses it.
  2795. if (modelProviderNamespaces().has(name)) queue.push(frame({ type: 'host/models-changed' }))
  2796. }),
  2797. ctx.on('credentials/updated', (ref) => {
  2798. queue.push(frame({ type: 'host/credentials-changed', ref: String(ref) }))
  2799. }),
  2800. ctx.on('llm/adapters-updated', () => {
  2801. queue.push(frame({ type: 'host/models-changed' }))
  2802. }),
  2803. ]
  2804. return queue.iterate(signal, () => { for (const dispose of disposers) dispose() })
  2805. },
  2806. },
  2807. respond(message: ClientResponse): Promise<RpcReceipt> {
  2808. // Route by the echoed rpcId (the wire correlation): approvals first,
  2809. // then questions — the two registries share one id space of UUIDs.
  2810. const approval = pendingApprovals.get(message.rpcId)
  2811. if (approval !== undefined) {
  2812. if (!message.result.ok) return Promise.resolve({ accepted: false, reason: 'bad-response' })
  2813. const parsed = approvalResponsePayloadSchema.safeParse(message.result.value)
  2814. // The payload's audit correlation must match the entry the rpcId routed
  2815. // to — a mismatched answer is malformed, not merely late.
  2816. if (!parsed.success || parsed.data.approvalId !== approval.approvalId || parsed.data.sessionId !== approval.sessionId) {
  2817. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  2818. }
  2819. approval.resolve(parsed.data.outcome)
  2820. return Promise.resolve({ accepted: true })
  2821. }
  2822. const pending = pendingQuestions.get(message.rpcId)
  2823. if (pending === undefined) return Promise.resolve({ accepted: false, reason: 'not-pending' })
  2824. if (!message.result.ok) {
  2825. if (message.result.error.code !== 'cancelled') {
  2826. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  2827. }
  2828. claimQuestion(pending, 'cancelled')
  2829. pending.reject(new UserInteractionError(
  2830. 'the user cancelled ask_user_question', 'ASK_CANCELLED'))
  2831. return Promise.resolve({ accepted: true })
  2832. }
  2833. const parsed = questionResponsePayloadSchema.safeParse(message.result.value)
  2834. if (!parsed.success) {
  2835. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  2836. }
  2837. const payload: QuestionResponsePayload = {
  2838. sessionId: parsed.data.sessionId,
  2839. answer: {
  2840. answers: parsed.data.answer.answers.map(answer => ({
  2841. id: answer.id,
  2842. selected: answer.selected,
  2843. ...(answer.custom === undefined ? {} : { custom: answer.custom }),
  2844. })),
  2845. },
  2846. }
  2847. if (!matchesQuestions(payload, pending)) {
  2848. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  2849. }
  2850. claimQuestion(pending, 'answered')
  2851. pending.resolve(payload.answer)
  2852. return Promise.resolve({ accepted: true })
  2853. },
  2854. }
  2855. }