api-proxy.ts 117 KB

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