api-proxy.ts 96 KB

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