| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501250225032504250525062507250825092510251125122513251425152516251725182519252025212522252325242525252625272528252925302531253225332534253525362537253825392540254125422543254425452546254725482549255025512552255325542555255625572558255925602561256225632564256525662567256825692570257125722573257425752576257725782579258025812582258325842585258625872588258925902591259225932594259525962597259825992600260126022603260426052606260726082609261026112612261326142615261626172618261926202621262226232624262526262627262826292630263126322633263426352636263726382639264026412642264326442645264626472648264926502651265226532654265526562657265826592660266126622663266426652666266726682669267026712672267326742675267626772678267926802681 |
- /**
- * Host-side ApiProxy implementation. Signature discipline: unary takes the
- * narrow RpcRequest<P> and echoes request.rpcId on the RpcResponse<T>.
- */
- import { randomUUID } from 'node:crypto'
- import { mkdir, stat } from 'node:fs/promises'
- import { join } from 'node:path'
- import type { Context } from 'cordis'
- import { installAgentLlmTarget } from '@deepseek-ai/dsh-agent'
- import type {
- Agent, AgentLlmTarget, AgentLlmTargetRef, AgentStatus, InboxItem, InboxItemId,
- } from '@deepseek-ai/dsh-agent'
- import { createUserMessage, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
- import { errorChain } from '@deepseek-ai/dsh-llm'
- import type { MessageSource } from '@deepseek-ai/dsh-llm'
- import { isAppendSurfaceEvent, lastActivityTime } from '@deepseek-ai/dsh-session'
- import type { Session, SessionEvent, SessionHeader, SessionId, UserMessage } from '@deepseek-ai/dsh-session'
- import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
- import { SessionQueryError, type SessionSearchCursor } from '@deepseek-ai/dsh-session-query'
- import { SubagentError } from '@deepseek-ai/dsh-subagent'
- import type { SubagentListEntry as CatalogSubagentListEntry } from '@deepseek-ai/dsh-subagent'
- import type { Workspace, WorkspaceRecord } from '@deepseek-ai/dsh-workspace'
- import {
- workspaceDomainState, workspaceRecord, WorkspaceId as brandWorkspaceId,
- WorkspaceMoveInvalidError, WorkspaceUnknownSessionError,
- } from '@deepseek-ai/dsh-workspace'
- // Type-only: brings the `ctx.tools` Context merge into this program (viewFor reads presenters).
- import type {} from '@deepseek-ai/dsh-tools'
- import type {
- ApiProxy, CredentialView, GoalRef, HistoryEntry, HostFrame, ModelCatalogFailure, ModelProviderGroup,
- ModelReasoning, MuxFrame, QuestionResponsePayload, SessionProjectionsBlock, SessionSearchItem,
- SessionSummary, SettingsNamespaceView, SubagentAddress, ToolEventView,
- WorkspaceId, WorkspaceView,
- } from './api/index.ts'
- import {
- SESSION_SEARCH_RESULT_LIMIT,
- SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS,
- truncateUnicodeCodePoints,
- } from './api/session-search.ts'
- // Type-only: resolves `ctx.get('sessionProjections')` to the projection registry.
- import type {} from '@deepseek-ai/dsh-session-projection'
- // Type-only: resolves `ctx.get('sessionProjectionCache')` (the cold listing column).
- import type {} from '@deepseek-ai/dsh-session-projection-cache'
- // GoalError narrows domain rejections to their stable codes at the wire boundary.
- import { GoalError } from '@deepseek-ai/dsh-goal'
- import type { GoalRef as CoreGoalRef } from '@deepseek-ai/dsh-goal'
- // Type-only edges: resolve `ctx.get('commands')`, the `commands/change` event, and `ctx.get('skills')`.
- import type {} from '@deepseek-ai/dsh-commands'
- import type {} from '@deepseek-ai/dsh-skill'
- // The settings/credentials seams: brand guards run at this wire boundary; the
- // service reads stay optional (`ctx.get`) so a composition without either
- // provider still serves every other domain.
- import { SettingsConflictError, settingsNamespace } from '@deepseek-ai/dsh-settings'
- import type { SettingsDescriptor, SettingsNamespace, SettingsPathOp } from '@deepseek-ai/dsh-settings'
- import { credentialRef } from '@deepseek-ai/dsh-credentials'
- // Value edge: the rename impl narrows the title service's validation failure; the import also resolves `ctx.get('sessionTitle')`.
- import { SessionTitleInvalidError } from '@deepseek-ai/dsh-session-title'
- import type { CallId } from '@deepseek-ai/dsh-llm/brand'
- import type { ApprovalOutcome, ApprovalRequestId } from '@deepseek-ai/dsh-user-approval'
- // Side-effect type import: resolves the `approval/request` waterfall and
- // `ctx.get('approval')` without a value dependency on the seam (optional composition).
- import type {} from '@deepseek-ai/dsh-user-approval'
- import { approvalResponsePayloadSchema } from './api/approvals.schema.ts'
- import { questionResponsePayloadSchema } from './api/questions.schema.ts'
- import type { ClientResponse, RpcError, RpcReceipt, RpcRequest, RpcResponse } from './api/rpc.ts'
- import { RpcId } from './api/rpc.ts'
- import type {
- AskUserQuestionAnswer, AskUserQuestionItem, AskUserQuestionRequest,
- } from '@deepseek-ai/dsh-user-interaction'
- import { UserInteractionError } from '@deepseek-ai/dsh-user-interaction'
- import { DirectoryPickerError } from '@deepseek-ai/dsh-host-directory-picker'
- import { openNativePath } from './native-path-opener.ts'
- /** Page size when history is called without maxMessages. */
- const DEFAULT_MAX_MESSAGES = 50
- /** Non-model settings namespaces intentionally served to the Web client. */
- const WEB_SETTINGS_NAMESPACES = ['permission'] as const
- /** Provider work budget: at most 100 calls and 2,000 inspected hits. */
- const SESSION_SEARCH_PROVIDER_CALL_LIMIT = 100
- /** Bound cold-log stat fan-out and settle each started batch before cancellation returns. */
- const COLD_SUMMARY_BATCH_SIZE = 16
- /** Conversation message event types (the pagination counting unit). */
- const MESSAGE_TYPES = new Set(['user/message', 'assistant/message', 'steering/message'])
- /** Product settings intentionally exposed beside model-provider namespaces. */
- const PRODUCT_SETTINGS_NAMESPACES = new Set(['ui-onboarding'])
- /** Read live abort state across awaits without treating it as synchronously immutable. */
- function isAborted(signal: AbortSignal): boolean {
- return signal.aborted
- }
- /**
- * Message-boundary pagination: count maxMessages append-origin messages
- * backwards from the window tail. Replacement copies never entered the
- * conversation a reader sees — they restate a shadowed range for the model
- * alone — so they consume no quota; the page stays one contiguous raw range,
- * which keeps a compaction's log-only provenance on the same page as its
- * replacement. The cut is the starting seq of the oldest message group (chunks
- * group via sourceEventSeqs — never cut mid-message). The tail page naturally
- * includes the in-progress partial.
- */
- function paginate(
- events: readonly SessionEvent[],
- beforeSeq: number | undefined,
- maxMessages: number,
- ): { events: SessionEvent[]; hasMore: boolean } {
- const window = beforeSeq === undefined ? [...events] : events.filter(event => event.seq < beforeSeq)
- let count = 0
- let cut = 0
- for (let i = window.length - 1; i >= 0; i--) {
- const event = window[i] as SessionEvent
- if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue
- count++
- const sources = (event as { sourceEventSeqs?: number[] }).sourceEventSeqs
- const groupStart = sources !== undefined && sources.length > 0 ? Math.min(event.seq, ...sources) : event.seq
- if (count >= maxMessages) {
- cut = groupStart
- break
- }
- }
- const page = window.filter(event => event.seq >= cut)
- return { events: page, hasMore: cut > 0 }
- }
- /** Wrap an ok result echoing the request's rpcId. */
- function ok<T>(request: RpcRequest<unknown>, value: T): RpcResponse<T> {
- return { rpcId: request.rpcId, result: { ok: true, value } }
- }
- /**
- * Build the provider/model catalog over every registered route. Shared by the
- * session-scoped `session.models` (which passes the session's current target
- * so an unlisted current model still renders selectable) and the host-scoped
- * `llm.models` (no current). Per-provider failures ride `failures` without
- * failing the sound groups; groups that advertise nothing are dropped.
- */
- async function buildModelCatalog(
- ctx: Context,
- current?: { provider: string; model: string },
- ): Promise<{ groups: ModelProviderGroup[]; failures: ModelCatalogFailure[] }> {
- const catalog = await Promise.all(ctx.llm.listProviders().map(async (provider) => {
- try {
- const advertised = await ctx.llm.listModels(provider.id)
- const models = [...advertised]
- if (
- current !== undefined
- && provider.id === current.provider
- && !models.some(model => model.id === current.model)
- ) {
- models.push({
- provider: provider.id,
- id: current.model,
- name: current.model,
- })
- }
- const entries = await Promise.all(models.map(async (model) => {
- const resolved = await ctx.llm.resolveModelInfo(provider.id, model.id)
- const reasoning: ModelReasoning | undefined = resolved.reasoning === undefined
- ? undefined
- : {
- efforts: resolved.reasoning.efforts.map(effort => ({
- id: effort.id,
- name: effort.name,
- ...effort.description === undefined
- ? {}
- : { description: effort.description },
- })),
- ...resolved.reasoning.defaultEffort === undefined
- ? {}
- : { defaultEffort: resolved.reasoning.defaultEffort },
- }
- return {
- id: model.id,
- name: model.name,
- ...model.description === undefined ? {} : { description: model.description },
- ...current !== undefined
- && provider.id === current.provider
- && model.id === current.model
- && !advertised.some(candidate => candidate.id === current.model)
- ? { unlisted: true as const }
- : {},
- ...reasoning === undefined ? {} : { reasoning },
- }
- }))
- const group: ModelProviderGroup = {
- id: provider.id,
- name: provider.name,
- models: entries,
- }
- return { kind: 'group' as const, group }
- } catch (error: unknown) {
- const failure: ModelCatalogFailure = {
- id: provider.id,
- name: provider.name,
- message: error instanceof Error ? error.message : String(error),
- }
- return { kind: 'failure' as const, failure }
- }
- }))
- return {
- groups: catalog.flatMap(item => item.kind === 'group' ? [item.group] : []).filter(group => group.models.length > 0),
- failures: catalog.flatMap(item => item.kind === 'failure' ? [item.failure] : []),
- }
- }
- /** Wrap an error result echoing the request's rpcId. */
- function err<T>(request: RpcRequest<unknown>, error: RpcError): RpcResponse<T> {
- return { rpcId: request.rpcId, result: { ok: false, error } }
- }
- /** Simple async queue: core callbacks push, the AsyncIterable pulls; abort/return cleans up. */
- class FrameQueue<F> {
- private buffer: F[] = []
- private waiter: (() => void) | undefined
- private done = false
- push(item: F): void {
- if (this.done) return
- this.buffer.push(item)
- this.waiter?.()
- }
- end(): void {
- this.done = true
- this.waiter?.()
- }
- async *iterate(signal: AbortSignal, cleanup: () => void): AsyncGenerator<F> {
- const onAbort = (): void => { this.end() }
- signal.addEventListener('abort', onAbort, { once: true })
- try {
- while (true) {
- while (this.buffer.length > 0) yield this.buffer.shift() as F
- if (this.done || signal.aborted) return
- await new Promise<void>((resolve) => { this.waiter = resolve })
- this.waiter = undefined
- }
- } finally {
- signal.removeEventListener('abort', onAbort)
- cleanup()
- }
- }
- }
- /**
- * Server-side frame mint: pure pushes get a fresh rpcId per frame (answerable
- * frames — approval/question requested — mint their stable id in their
- * pending registries instead).
- */
- function frame<F>(payload: F): RpcRequest<F> {
- return { rpcId: RpcId(randomUUID()), payload }
- }
- /** Queue the subscription baseline frame. */
- function subscribeSession(queue: FrameQueue<RpcRequest<MuxFrame>>, session: Session): void {
- queue.push(frame({ type: 'session/subscribed', sessionId: session.id, lastSeq: session.seq - 1 }))
- }
- /**
- * Whether the session's conversation has started: no turn has run yet (a
- * turn is one model-loop execution). Standalone plugin events — command
- * lifecycle records, plan/mode, titles, goals — never open a turn, so
- * running `/plan` or `/goal` on a fresh session keeps it blank
- * (list-hidden, reusable).
- */
- function sessionBlank(session: Session): boolean {
- return !session.events.some(event => event.type === 'turn/start')
- }
- /** Shared Session-header projection for list baselines and creation frames. */
- function sessionListFields(header: SessionHeader): {
- parentSessionId?: SessionId
- origin?: 'subagent'
- cwd?: string
- } {
- return {
- ...header.parentSession === undefined ? {} : { parentSessionId: header.parentSession },
- ...header.origin === undefined ? {} : { origin: header.origin },
- ...header.cwd === undefined ? {} : { cwd: header.cwd },
- }
- }
- /** SessionSummary projection for attached (in-memory) sessions. */
- function summarize(session: Session, running: boolean): SessionSummary {
- return {
- sessionId: session.id,
- // Excludes end-seed: a resumed-but-untouched session
- // must not sort as freshly worked in.
- updatedAt: lastActivityTime(session.events) ?? session.header.createdAt,
- running,
- blank: sessionBlank(session),
- ...sessionListFields(session.header),
- }
- }
- /**
- * SessionSummary projection for cold (persisted, unattached) sessions.
- * updatedAt is the log file's mtime; backends without a per-session file
- * (locate() undefined) fall back to the header's createdAt.
- */
- async function summarizeCold(
- persistence: SessionPersistence,
- meta: SessionHeader,
- signal?: AbortSignal,
- ): Promise<SessionSummary> {
- signal?.throwIfAborted()
- let updatedAt = meta.createdAt
- const location = persistence.locate(meta)
- signal?.throwIfAborted()
- if (location !== undefined) {
- try {
- updatedAt = (await stat(location.path)).mtimeMs
- } catch {
- // The log vanished between list() and stat() (concurrent cleanup); createdAt stands in.
- }
- signal?.throwIfAborted()
- }
- return {
- sessionId: meta.id,
- updatedAt,
- running: false,
- // Lazy persistence keeps never-appended sessions out of list(); reading
- // a cold log to check for turns would defeat the index read, so a listed
- // cold session is served as not-blank (its log holds its conversation).
- blank: false,
- ...meta.parentSession === undefined ? {} : { parentSessionId: meta.parentSession },
- ...meta.origin === undefined ? {} : { origin: meta.origin },
- /* v8 ignore next -- the empty arm needs a cwd-less meta, but list()
- filters those out (legacy logs are not served); the conditional mirrors
- summarize() shape. */
- ...meta.cwd === undefined ? {} : { cwd: meta.cwd },
- }
- }
- /** Map a browse-primitive failure onto the wire error vocabulary (unknown throws stay internal). */
- function directoryError(error: unknown): RpcError {
- if (error instanceof DirectoryPickerError) {
- return { code: error.code, message: error.message, details: { path: error.path } }
- }
- return { code: 'internal', message: error instanceof Error ? error.message : String(error), details: {} }
- }
- /** Resolved Host routing and project-directory defaults consumed by the API implementation. */
- export interface ApiProxyDefaults {
- provider: string
- model: string
- /** Default project directory for new sessions whose create request carries no cwd. */
- cwd: string
- /** Parent directory for name-created workspaces. */
- workspaceRoot: string
- /** Native open-with-default-application; injectable for carrier tests. */
- openPath?: (path: string, signal: AbortSignal) => Promise<void>
- }
- /** The tool/call payload fields the presenter path reads. */
- interface ToolCallData { callId: string; name: string; arguments: string }
- /**
- * One outstanding approval question: the stable server-request id, the frame
- * material replayed to late mux subscribers, and the resolver that settles the
- * answerer's promise back into `ctx.approval`.
- */
- interface PendingApproval {
- rpcId: RpcId
- sessionId: SessionId
- approvalId: ApprovalRequestId
- toolName: string
- callId?: CallId
- reason?: string
- resolve(outcome: ApprovalOutcome): void
- }
- /** Project a pending entry into its answerable mux frame (initial push and mux-open replay share it). */
- function requestedFrame(pending: PendingApproval): RpcRequest<MuxFrame> {
- return {
- rpcId: pending.rpcId,
- payload: {
- type: 'approval/requested',
- sessionId: pending.sessionId,
- approvalId: pending.approvalId,
- toolName: pending.toolName,
- ...pending.callId === undefined ? {} : { callId: pending.callId },
- ...pending.reason === undefined ? {} : { reason: pending.reason },
- },
- }
- }
- /** One host-owned question wait, addressed by the stable server-request id. */
- interface PendingQuestion {
- rpcId: RpcId
- sessionId: SessionId
- questions: AskUserQuestionItem[]
- resolve: (answer: AskUserQuestionAnswer) => void
- reject: (error: UserInteractionError) => void
- signal?: AbortSignal
- onAbort?: () => void
- }
- /** Validate one answer batch against the exact question request it resolves. */
- function matchesQuestions(payload: QuestionResponsePayload, pending: PendingQuestion): boolean {
- if (payload.sessionId !== pending.sessionId) return false
- const answers = payload.answer.answers
- if (answers.length !== pending.questions.length) return false
- return answers.every((answer, index) => {
- const question = pending.questions[index] as AskUserQuestionItem
- if (answer.id !== question.id) return false
- if (new Set(answer.selected).size !== answer.selected.length) return false
- const custom = answer.custom?.trim()
- if (custom !== undefined && custom === '') return false
- if (custom !== undefined && answer.selected.length > 0) return false
- if (question.multiSelect !== true && answer.selected.length > 1) return false
- const labels = new Set(question.options?.map(option => option.label) ?? [])
- return answer.selected.every(label => labels.has(label))
- })
- }
- /**
- * Compute the render intent for a tool/call or tool/result event through the
- * presenters registered at this moment; every other event type gets none. A
- * result's presenter needs its call's parsed args — `argsFor` supplies them
- * (live: the per-session call table; history: an in-page backscan), returning
- * undefined when the pairing is unavailable (e.g. the call fell off the page),
- * which soft-falls to no view. Presenter or JSON.parse throws also soft-fall:
- * the client's documented default (generic JSON card) covers every miss.
- */
- function viewFor(ctx: Context, event: SessionEvent, argsFor: (callId: string) => unknown): ToolEventView | undefined {
- try {
- if (event.type === 'tool/call') {
- const { name, arguments: raw } = event.data as ToolCallData
- const view = ctx.tools.get(name)?.presentCall?.(JSON.parse(raw))
- return view === undefined ? undefined : { for: 'call', view }
- }
- if (event.type === 'tool/result') {
- const { message, meta } = event.data
- const [result] = message.content
- const callId = message.source.callId
- const call = argsFor(callId) as { name: string; args: unknown } | undefined
- if (call === undefined) return undefined
- const view = ctx.tools.get(call.name)?.presentResult?.(call.args, {
- content: result.content,
- isError: result.isError === true,
- ...meta === undefined ? {} : { meta },
- })
- return view === undefined ? undefined : { for: 'result', view }
- }
- } catch (error: unknown) {
- // A throwing presenter (or unparseable arguments) must not break delivery;
- // the event still ships, just without a view.
- console.error(`api-proxy: presenter failed for ${event.type}, falling back to generic: ${String(error)}`)
- }
- return undefined
- }
- /**
- * Resolve a tool/result's call pairing by scanning a window of events backwards
- * for the matching tool/call. Used by the history path (the page is the
- * window — a cross-page pairing soft-falls to no view) and by live-path table
- * misses after a reconnect-eviction.
- */
- function backscanArgs(events: readonly SessionEvent[], callId: string): { name: string; args: unknown } | undefined {
- for (let i = events.length - 1; i >= 0; i--) {
- const event = events[i] as SessionEvent
- if (event.type !== 'tool/call') continue
- const data = event.data as ToolCallData
- if (data.callId !== callId) continue
- try {
- return { name: data.name, args: JSON.parse(data.arguments) }
- } catch {
- // Unparseable stored arguments: same soft-fall as a live parse failure.
- return undefined
- }
- }
- return undefined
- }
- /** Render one detached history page through the same presenter path as ordinary history. */
- function historyPage(
- ctx: Context,
- events: readonly SessionEvent[],
- beforeSeq: number | undefined,
- maxMessages: number | undefined,
- ): { events: HistoryEntry[]; hasMore: boolean } {
- const page = paginate(events, beforeSeq, maxMessages ?? DEFAULT_MAX_MESSAGES)
- return {
- events: page.events.map((event) => {
- const view = viewFor(ctx, event, callId => backscanArgs(page.events, callId))
- return { event, ...view === undefined ? {} : { view } }
- }),
- hasMore: page.hasMore,
- }
- }
- /**
- * The projection baseline for one history tail page: the registry's
- * watermark-cache snapshot — one fully synchronous read (no await between the
- * page slice and this), so all values and `asOfSeq` form a single consistent
- * cut and `asOfSeq` equals the window tail event seq. The carrier holds zero
- * domain knowledge (each value passed its unit's own schema inside the
- * registry). An absent registry means the deployment has no projection seam:
- * the whole block is absent and clients treat every key as capability-absent.
- */
- function projectionsFor(ctx: Context, session: Session): SessionProjectionsBlock | undefined {
- const registry = ctx.get('sessionProjections')
- if (registry === undefined) return undefined
- return registry.snapshot(session)
- }
- /**
- * The projection baseline of one session.list row, fail-soft: attached
- * sessions cut the registry's live watermark cache; cold sessions view the
- * persisted projection cache's identity-checked stored rows (zero log loads
- * either way — the listing use case the cache exists for). The block shape
- * (values + asOfSeq) matches the history tail's, so a client seeds its
- * value store under the same higher-seq-wins rule. Any failure — and an
- * empty value set — yields an absent block: a listing without projections
- * is degraded, never broken.
- */
- function listProjectionsFor(ctx: Context, meta: SessionHeader, session: Session | undefined): SessionProjectionsBlock | undefined {
- try {
- const block = session !== undefined
- ? ctx.get('sessionProjections')?.snapshot(session)
- : ctx.get('sessionProjectionCache')?.cachedSnapshot(meta)
- return block !== undefined && Object.keys(block.values).length > 0 ? block : undefined
- } catch (error) {
- ctx.logger.warn(`session.list: projection column for "${meta.id}" failed (serving the row without it): ${String(error)}`)
- return undefined
- }
- }
- /** Projection baseline for a detached history tail without Agent activation. */
- function detachedProjectionsFor(
- ctx: Context,
- events: readonly SessionEvent[],
- ): SessionProjectionsBlock | undefined {
- const registry = ctx.get('sessionProjections')
- if (registry === undefined) return undefined
- return registry.restore({}, events, 0).snapshot
- }
- /** Map continuation admission failures without exposing provider details. */
- function subagentPromptError(
- request: RpcRequest<{ childSessionId: SessionId }>,
- error: unknown,
- signal: AbortSignal,
- ): RpcResponse<never> {
- const childSessionId = request.payload.childSessionId
- if (signal.aborted) {
- return err(request, { code: 'cancelled', message: 'subagent prompt was cancelled', details: {} })
- }
- if (error instanceof SubagentError) {
- switch (error.code) {
- case 'NOT_RESUMABLE':
- return err(request, {
- code: 'subagent-not-resumable',
- message: 'subagent cannot be resumed',
- details: { childSessionId },
- })
- case 'UNAUTHORIZED':
- return err(request, {
- code: 'subagent-unauthorized',
- message: 'subagent does not belong to this parent',
- details: { childSessionId },
- })
- case 'DRAINING':
- case 'ACTIVATION_CLOSING':
- case 'CONTINUATION_UNAVAILABLE':
- case 'PERSISTENCE_UNAVAILABLE':
- return err(request, {
- code: 'subagent-delivery-unavailable',
- message: 'subagent follow-up is temporarily unavailable',
- details: { childSessionId },
- })
- default:
- break
- }
- }
- return err(request, { code: 'internal', message: 'subagent prompt failed', details: {} })
- }
- /** Verify one address and mode against the complete direct-child catalog. */
- async function catalogChild(
- ctx: Context,
- address: SubagentAddress,
- signal?: AbortSignal,
- ): Promise<{
- entry?: Extract<CatalogSubagentListEntry, { kind: 'child' }>
- error?: RpcError
- }> {
- const { parentSessionId, childSessionId, mode } = address
- try {
- const entries = await ctx.subagents.listChildren(parentSessionId, signal)
- const entry = entries.find(candidate => candidate.id === childSessionId)
- if (entry === undefined || (entry.kind === 'child' && entry.mode !== mode)) {
- return {
- error: {
- code: 'subagent-not-found',
- message: `session "${childSessionId}" is not a ${mode} direct child of "${parentSessionId}"`,
- details: { parentSessionId, childSessionId },
- },
- }
- }
- if (entry.kind === 'diagnostic') {
- return {
- error: {
- code: 'subagent-catalog-diagnostic',
- message: `subagent "${childSessionId}" is ${entry.reason}`,
- details: { parentSessionId, childSessionId, reason: entry.reason },
- },
- }
- }
- return { entry }
- } catch (error: unknown) {
- if (signal?.aborted
- || (error instanceof SubagentError && error.code === 'CANCELLED')
- || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')) {
- return { error: { code: 'cancelled', message: 'subagent catalog read was cancelled', details: {} } }
- }
- if (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') {
- return {
- error: {
- code: 'subagent-not-found',
- message: `parent session "${parentSessionId}" was not found`,
- details: { parentSessionId, childSessionId },
- },
- }
- }
- return { error: { code: 'internal', message: 'subagent catalog read failed', details: {} } }
- }
- }
- /**
- * Thrown by the cold-resume path when the id names no servable session
- * (absent from the store, or a pre-project legacy log without a cwd).
- */
- class SessionNotFound extends Error {}
- /** Session identity whose lifecycle belongs to subagent routing, not generic Host resume. */
- class SubagentSessionOwnership extends Error {
- constructor(readonly sessionId: SessionId) {
- super(`session "${sessionId}" is a subagent session; use subagent delivery`)
- }
- }
- /** Requested identity already belongs to a session with another project cwd. */
- class SessionCwdConflict extends Error {
- constructor(
- readonly sessionId: SessionId,
- readonly requestedCwd: string,
- readonly existingCwd: string | undefined,
- ) {
- super(
- `session "${sessionId}" already exists with cwd ${JSON.stringify(existingCwd)}; `
- + `requested ${JSON.stringify(requestedCwd)}`,
- )
- }
- }
- /** Host failed before the registry could adopt a name-created directory. */
- class WorkspaceDirectoryCreationError extends Error {}
- /** An explicit Host naming operation would duplicate another Workspace title. */
- class WorkspaceNameConflictError extends Error {
- constructor(readonly workspaceName: string) {
- super(`workspace name '${workspaceName}' is already in use`)
- this.name = 'WorkspaceNameConflictError'
- }
- }
- /** Shared workspace-not-found error response of the workspace.* mutation rows. */
- function workspaceNotFound<T>(request: RpcRequest<unknown>, workspaceId: string): RpcResponse<T> {
- return err(request, {
- code: 'workspace-not-found',
- message: `workspace "${workspaceId}" not found`,
- details: { workspaceId },
- })
- }
- /** Wire projection of one workspace entity (the workspace.* value row). */
- function workspaceView(workspace: Workspace): WorkspaceView {
- return {
- workspaceId: workspace.id,
- path: workspace.path,
- title: workspace.title,
- sessionIds: [...workspace.sessionIds],
- createdAt: workspace.createdAt,
- updatedAt: workspace.updatedAt,
- }
- }
- /** Wire projection of the durable record carried by `domain/changed`. */
- function changedWorkspaceView(workspaceId: string, value: unknown): WorkspaceView {
- const record: WorkspaceRecord = workspaceRecord.parse(value)
- return {
- workspaceId: workspaceId as WorkspaceId,
- path: record.path,
- title: record.title,
- sessionIds: [...record.sessionIds],
- createdAt: record.createdAt,
- updatedAt: record.updatedAt,
- }
- }
- /**
- * Implement ApiProxy over a composed host context.
- * @param ctx - a context with the Host spine and Workspace registry mounted.
- * @param defaults - host routing and project-directory defaults.
- * @returns the ApiProxy implementation.
- */
- export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiProxy {
- const agentOptions = { provider: defaults.provider, model: defaults.model }
- type WebLlmTargetRef = AgentLlmTargetRef & { current: AgentLlmTarget }
- const targets = new WeakMap<Agent, WebLlmTargetRef>()
- /** Implicit resume of cold sessions, deduplicating concurrent calls (follows the jsonrpc sessionCreations precedent). */
- const resumes = new Map<SessionId, Promise<Agent>>()
- /** Client-chosen identity creation/resume, deduplicated across concurrent retries. */
- const sessionCreations = new Map<SessionId, Promise<Agent>>()
- /** Serializes path ownership and explicit title checks with Workspace mutations. */
- let workspaceCreationChain = Promise.resolve()
- const pendingQuestions = new Map<RpcId, PendingQuestion>()
- const pendingApprovals = new Map<RpcId, PendingApproval>()
- const muxQueues = new Set<FrameQueue<RpcRequest<MuxFrame>>>()
- /**
- * Install or return the session-local target that prompt assembly snapshots.
- * Seed order: latest logged request/header, else the host default routing.
- * There is no create-time per-session override tier on this wire — if one
- * returns (a create-options contribution), it must fold in between the two.
- */
- function targetFor(agent: Agent): WebLlmTargetRef {
- const installed = targets.get(agent)
- if (installed !== undefined) return installed
- const logged = agent.session.requestHeader()?.config
- const target: WebLlmTargetRef = {
- current: logged === undefined
- ? { provider: defaults.provider, model: defaults.model }
- : {
- provider: logged.provider,
- model: logged.model,
- ...logged.reasoningEffort === undefined
- ? {}
- : { reasoningEffort: logged.reasoningEffort },
- },
- assembled: undefined,
- }
- installAgentLlmTarget(agent.ctx, target)
- targets.set(agent, target)
- return target
- }
- /** Pre-publication setup used by both fresh and resumed Web agents. */
- function installTarget(agentCtx: Context): void {
- const agent = agentCtx.agent
- if (agent === undefined) throw new Error('api-proxy: agent setup has no scoped agent')
- targetFor(agent)
- }
- /** Send one transient frame to every connected mux consumer. */
- function broadcast(payload: MuxFrame): void {
- const envelope = frame(payload)
- for (const queue of muxQueues) queue.push(envelope)
- }
- // Projection change feed → session/projection push frames. The carrier
- // mints the wire frame (the seam package holds no wire vocabulary); the
- // child activates only when a projection registry is composed, and the
- // subscription unwinds with this gateway's fiber.
- ctx.inject(['sessionProjections'], (projectionCtx) => {
- projectionCtx.sessionProjections.onChanged((session, key, value, seq) => {
- broadcast({ type: 'session/projection', sessionId: session.id, key, value, seq })
- })
- })
- /**
- * Per-session queued-occurrence mirror serving the mux-open queue snapshot
- * (the same refresh-recovery baseline as pending questions). Each terminal
- * queue event retires one matching occurrence, so repeated sends of the same
- * identified message remain visible until every occurrence is claimed.
- */
- const queuedMirror = new Map<SessionId, InboxItem[]>()
- type UnseenQueueEvent =
- | { readonly kind: 'update'; readonly item: InboxItem }
- | { readonly kind: 'terminal' }
- const unseenQueueEvents = new Map<SessionId, Map<InboxItemId, UnseenQueueEvent>>()
- const rememberUnseen = (sessionId: SessionId, itemId: InboxItemId, event: UnseenQueueEvent): void => {
- let events = unseenQueueEvents.get(sessionId)
- if (events === undefined) {
- events = new Map()
- unseenQueueEvents.set(sessionId, events)
- }
- events.set(itemId, event)
- // Only synchronous re-entrancy may deliver a mutation before its outer
- // enqueue observer. Drop unmatched protocol-invalid observations instead
- // of retaining process-local ids indefinitely.
- queueMicrotask(() => {
- const current = unseenQueueEvents.get(sessionId)
- if (current?.get(itemId) !== event) return
- current.delete(itemId)
- if (current.size === 0) unseenQueueEvents.delete(sessionId)
- })
- }
- const takeUnseen = (sessionId: SessionId, itemId: InboxItemId): UnseenQueueEvent | undefined => {
- const events = unseenQueueEvents.get(sessionId)
- const event = events?.get(itemId)
- if (event === undefined) return undefined
- events?.delete(itemId)
- if (events?.size === 0) unseenQueueEvents.delete(sessionId)
- return event
- }
- const publishQueue = (sessionId: SessionId): void => {
- const items = queuedMirror.get(sessionId) ?? []
- broadcast({
- type: 'session/queue',
- sessionId,
- items: items.map(item => ({
- id: item.id,
- message: item.message,
- })),
- })
- }
- ctx.effect(() => {
- const retire = (agent: Agent, item: InboxItem): boolean => {
- const entries = queuedMirror.get(agent.id)
- if (entries === undefined) {
- rememberUnseen(agent.id, item.id, { kind: 'terminal' })
- return false
- }
- const index = entries.findIndex(entry => entry.id === item.id)
- if (index === -1) {
- rememberUnseen(agent.id, item.id, { kind: 'terminal' })
- return false
- }
- entries.splice(index, 1)
- if (entries.length === 0) queuedMirror.delete(agent.id)
- return true
- }
- const disposers = [
- ctx.on('agent/inbox/enqueue', (agent: Agent, item: InboxItem) => {
- if (item.placement !== 'queued') return
- const unseen = takeUnseen(agent.id, item.id)
- if (unseen?.kind === 'terminal') return
- let entries = queuedMirror.get(agent.id)
- if (entries === undefined) {
- entries = []
- queuedMirror.set(agent.id, entries)
- }
- entries.push(unseen?.kind === 'update' ? unseen.item : item)
- publishQueue(agent.id)
- }),
- ctx.on('agent/inbox/update', (agent: Agent, item: InboxItem) => {
- const entries = queuedMirror.get(agent.id)
- if (entries === undefined) {
- rememberUnseen(agent.id, item.id, { kind: 'update', item })
- return
- }
- const index = entries.findIndex(entry => entry.id === item.id)
- if (index === -1) {
- rememberUnseen(agent.id, item.id, { kind: 'update', item })
- return
- }
- entries.splice(index, 1, item)
- publishQueue(agent.id)
- }),
- ctx.on('agent/inbox/dequeue', (agent: Agent, item: InboxItem) => {
- if (retire(agent, item)) publishQueue(agent.id)
- }),
- ctx.on('agent/inbox/discard', (agent: Agent, items: InboxItem[]) => {
- let changed = false
- for (const item of items) changed = retire(agent, item) || changed
- if (changed) publishQueue(agent.id)
- }),
- ctx.on('session/disposed', (session: Session) => {
- queuedMirror.delete(session.id)
- unseenQueueEvents.delete(session.id)
- }),
- ]
- return () => { for (const dispose of disposers) dispose() }
- }, 'api-proxy: queued mirror')
- /** Remove a wait before settling it: synchronous deletion makes the first claimant win. */
- function claimQuestion(pending: PendingQuestion, outcome: 'answered' | 'cancelled'): void {
- pendingQuestions.delete(pending.rpcId)
- if (pending.signal !== undefined && pending.onAbort !== undefined) {
- pending.signal.removeEventListener('abort', pending.onAbort)
- }
- broadcast({
- type: 'question/resolved', sessionId: pending.sessionId,
- questionRpcId: pending.rpcId, outcome,
- })
- }
- const disposeProvider = ctx.userInteraction.registerProvider({
- ask(request: AskUserQuestionRequest): Promise<AskUserQuestionAnswer> {
- const sessionId = request.agent?.id
- if (sessionId === undefined) {
- return Promise.reject(new UserInteractionError(
- 'web user interaction requires an agent-owned session', 'ASK_MISSING_AGENT'))
- }
- return new Promise<AskUserQuestionAnswer>((resolve, reject) => {
- const rpcId = RpcId(randomUUID())
- const pending: PendingQuestion = {
- rpcId, sessionId, questions: request.questions, resolve, reject,
- ...(request.signal === undefined ? {} : { signal: request.signal }),
- }
- const onAbort = (): void => {
- claimQuestion(pending, 'cancelled')
- reject(new UserInteractionError(
- 'ask_user_question was aborted before the user answered', 'ASK_ABORTED'))
- }
- pending.onAbort = onAbort
- pendingQuestions.set(rpcId, pending)
- request.signal?.addEventListener('abort', onAbort, { once: true })
- const envelope: RpcRequest<MuxFrame> = {
- rpcId,
- payload: { type: 'question/requested', sessionId, questions: request.questions },
- }
- for (const queue of muxQueues) queue.push(envelope)
- })
- },
- })
- ctx.effect(() => () => {
- disposeProvider()
- for (const pending of [...pendingQuestions.values()]) {
- claimQuestion(pending, 'cancelled')
- pending.reject(new UserInteractionError(
- 'web user-interaction provider was disposed', 'ASK_ABORTED'))
- }
- }, 'api-proxy: user-interaction provider')
- // --- Approval pending registry ------------------------------------------
- // The proxy is the approval channel for every agent this host owns: an ask
- // through `ctx.approval` becomes an answerable server-request on the mux
- // stream (stable rpcId), settled by POST /api/respond. The entry survives
- // client disconnects — mux-open replays still-pending requested frames with
- // the same rpcId (the refresh-recovery baseline) — and withdraws on the
- // ask's own abort signal (turn cancel), pushing `cancelled` to subscribers.
- if (ctx.get('approval') !== undefined) {
- // Teardown parity with the question provider above: a gateway disposed
- // while approvals are pending settles every entry as 'cancelled' (the
- // service's fail-closed vocabulary), so no ask promise dangles past the
- // proxy's lifetime and subscribers see the withdrawal.
- ctx.effect(() => () => {
- for (const pending of [...pendingApprovals.values()]) pending.resolve('cancelled')
- }, 'api-proxy: approval registry teardown')
- ctx.on('approval/request', (req, next) => {
- // Dispatch rides a microtask behind the service's own signal check: an
- // abort landing in that window would register the abort listener AFTER
- // the signal fired — never invoked, entry pending forever, zombie frame
- // on every mux replay. Settle synchronously instead of publishing.
- if (req.signal?.aborted === true) return Promise.resolve<ApprovalOutcome>('cancelled')
- // The audit pair `approval/asked` is already appended by the service
- // before dispatch, but dispatch rides a microtask: parallel tool calls
- // can append several asked events before any answerer runs. THIS
- // request's event is therefore the newest asked event that is still
- // undecided, unclaimed by another pending entry, and — when the ask
- // names a call — carries the same callId.
- const events = req.agent.session.events
- const claimed = new Set<ApprovalRequestId>()
- for (const entry of pendingApprovals.values()) claimed.add(entry.approvalId)
- const decided = new Set<ApprovalRequestId>()
- let approvalId: ApprovalRequestId | undefined
- for (let i = events.length - 1; i >= 0; i -= 1) {
- const event = events[i] as SessionEvent
- if (event.type === 'approval/decided') {
- decided.add(event.data.id)
- } else if (event.type === 'approval/asked') {
- if (decided.has(event.data.id) || claimed.has(event.data.id)) continue
- // Symmetric pairing: a callId-bearing ask only takes its own call's
- // record, and a callId-less ask only takes a callId-less record —
- // so neither shape can steal the other's audit id under parallel
- // asks. (Today every producer — the tool executor — passes callId;
- // the callId-less arm guards any future non-tool asker.)
- if ((req.callId ?? null) !== (event.data.callId ?? null)) continue
- approvalId = event.data.id
- break
- }
- }
- // No asked event means the request bypassed the service's audit path —
- // not this channel's question; delegate to the fail-closed default.
- if (approvalId === undefined) return next()
- const id = approvalId
- return new Promise<ApprovalOutcome>((resolve) => {
- const settle = (outcome: ApprovalOutcome): void => {
- /* v8 ignore next 3 -- defensive double-settle guard: respond() routes
- through the pending table (a settled id is not-pending before it can
- re-settle) and the first settle removes the abort listener, so no
- reachable path settles twice; kept against future settle callers. */
- if (!pendingApprovals.delete(pending.rpcId)) return
- req.signal?.removeEventListener('abort', onAbort)
- broadcast({ type: 'approval/resolved', sessionId: pending.sessionId, approvalId: id, outcome })
- // A cancelled ask was already settled by the service's own signal
- // race, which discards this late resolution; resolving is a no-op
- // there and keeps this promise from dangling forever.
- resolve(outcome)
- }
- const onAbort = (): void => { settle('cancelled') }
- const pending: PendingApproval = {
- rpcId: RpcId(randomUUID()),
- sessionId: req.agent.session.id,
- approvalId: id,
- toolName: req.toolName,
- ...req.callId === undefined ? {} : { callId: req.callId },
- ...req.reason === undefined ? {} : { reason: req.reason },
- resolve: settle,
- }
- pendingApprovals.set(pending.rpcId, pending)
- req.signal?.addEventListener('abort', onAbort, { once: true })
- const envelope = requestedFrame(pending)
- for (const queue of muxQueues) queue.push(envelope)
- })
- })
- }
- /** Whether the session's own suffix carries the durable subagent discriminator. */
- function hasSubagentDescriptor(session: Pick<Session, 'events' | 'header'>): boolean {
- const ownStart = session.header.seedLength ?? 0
- return session.events.slice(ownStart).some(event => event.type === 'subagent/descriptor')
- }
- /**
- * Generic Host interaction cannot claim a durably classified subagent or an
- * Agent created through its live parent. The runtime-owner arm also covers
- * descriptor-less child publication windows and older stored headers.
- */
- function hasSubagentOwner(
- session: Pick<Session, 'events' | 'header'>,
- agent: Agent | undefined,
- ): boolean {
- if (session.header.origin === 'subagent' || hasSubagentDescriptor(session)) return true
- const parentId = session.header.parentSession
- if (parentId === undefined || agent === undefined) return false
- const parent = ctx.agents.get(parentId)
- return parent !== undefined && ctx.agents.isOwnedBy(agent.id, parent)
- }
- /** Stable generic-Host error for an identity reserved to subagent routing. */
- function subagentOwnershipError(sessionId: SessionId): RpcError {
- return {
- code: 'agent-busy',
- message: `session "${sessionId}" is owned by subagent routing`,
- details: { reason: 'use subagent delivery for this child session' },
- }
- }
- /** Inspect one cold served session without repairing, resuming, or publishing it. */
- async function inspectServable(sessionId: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
- const persistence = ctx.get('sessionPersistence')
- if (persistence === undefined) {
- throw new Error('session persistence is not configured (load a dsh-session-persistence backend)')
- }
- const meta = (await persistence.list()).find(m => m.id === sessionId)
- if (meta === undefined || meta.cwd === undefined) throw new SessionNotFound(`session "${sessionId}" not found`)
- const inspected = await persistence.inspect(sessionId)
- if (inspected.meta.cwd === undefined) throw new SessionNotFound(`session "${sessionId}" not found`)
- return inspected
- }
- async function agentFor(sessionId: SessionId): Promise<{ agent: Agent } | { error: RpcError }> {
- const attached = ctx.sessions.get(sessionId)
- const live = ctx.agents.get(sessionId)
- if (attached !== undefined && hasSubagentOwner(attached, live)) {
- return { error: subagentOwnershipError(sessionId) }
- }
- if (live !== undefined) return { agent: live }
- let resume = resumes.get(sessionId)
- if (resume === undefined) {
- resume = (async () => {
- try {
- const inspected = await inspectServable(sessionId)
- if (hasSubagentOwner({ header: inspected.meta, events: inspected.events }, undefined)) {
- throw new SubagentSessionOwnership(sessionId)
- }
- const publishedSession = ctx.sessions.get(sessionId)
- const publishedAgent = ctx.agents.get(sessionId)
- if (publishedSession !== undefined && hasSubagentOwner(publishedSession, publishedAgent)) {
- throw new SubagentSessionOwnership(sessionId)
- }
- const handle = await ctx.agents.resume({
- resumeSessionId: sessionId,
- agentOptions,
- setup: installTarget,
- })
- return handle.agent
- } finally {
- resumes.delete(sessionId)
- }
- })()
- resumes.set(sessionId, resume)
- }
- try {
- return { agent: await resume }
- } catch (error: unknown) {
- if (error instanceof SessionNotFound) {
- return { error: { code: 'session-not-found', message: error.message, details: { sessionId } } }
- }
- if (error instanceof SubagentSessionOwnership) {
- return { error: subagentOwnershipError(error.sessionId) }
- }
- // The internal details slot is contractually {}; the reason rides the message.
- return { error: { code: 'internal', message: `resume failed for session "${sessionId}": ${String(error)}`, details: {} } }
- }
- }
- type SessionReadState = {
- id: SessionId
- header: SessionHeader
- events: SessionEvent[]
- }
- /** Read one stable session prefix without acquiring an Agent owner. */
- async function readSessionState(sessionId: SessionId): Promise<SessionReadState> {
- const attached = ctx.sessions.get(sessionId)
- if (attached !== undefined) {
- return {
- id: attached.id,
- header: attached.header,
- events: [...attached.events],
- }
- }
- const inspected = await inspectServable(sessionId)
- return { id: inspected.meta.id, header: inspected.meta, events: inspected.events }
- }
- /** Resolve the Workspace inherited by a fork without making ordinary loose lineage grouped. */
- async function forkWorkspace(source: Pick<Session, 'id' | 'header'>): Promise<Workspace | undefined> {
- const workspaces = ctx.workspace.list()
- const direct = workspaces.find(workspace => workspace.sessionIds.includes(source.id))
- if (direct !== undefined || source.header.origin !== 'subagent') return direct
- const lineage = await ctx.sessionQuery.traceSession(source.id)
- for (const ancestor of lineage.ancestors) {
- const workspace = workspaces.find(candidate => candidate.sessionIds.includes(ancestor.header.id))
- if (workspace !== undefined) return workspace
- }
- return undefined
- }
- /** Read one transcript cut and optional projection baseline without acquiring an Agent owner. */
- async function historyStateFor(
- sessionId: SessionId,
- includeProjections: boolean,
- ): Promise<{ events: SessionEvent[]; projections?: SessionProjectionsBlock }> {
- const attached = ctx.sessions.get(sessionId)
- if (attached !== undefined) {
- const events = [...attached.events]
- const projections = includeProjections ? projectionsFor(ctx, attached) : undefined
- return { events, ...projections === undefined ? {} : { projections } }
- }
- const inspected = await inspectServable(sessionId)
- const projections = includeProjections ? detachedProjectionsFor(ctx, inspected.events) : undefined
- return {
- events: inspected.events,
- ...projections === undefined ? {} : { projections },
- }
- }
- /** Resolve one requested identity to a live agent, creating or resuming it once. */
- async function ensureSession(sessionId: SessionId, cwd: string, checkPersistedIdentity: boolean): Promise<Agent> {
- let creation = sessionCreations.get(sessionId)
- if (creation === undefined) {
- creation = (async () => {
- const attached = ctx.sessions.get(sessionId)
- const live = ctx.agents.get(sessionId)
- if (attached !== undefined && hasSubagentOwner(attached, live)) {
- throw new SubagentSessionOwnership(sessionId)
- }
- if (live !== undefined) return live
- const persistence = checkPersistedIdentity ? ctx.get('sessionPersistence') : undefined
- const stored = persistence === undefined
- ? undefined
- : (await persistence.list()).find(header => header.id === sessionId)
- if (persistence !== undefined && stored !== undefined) {
- if (stored.cwd !== cwd) {
- throw new SessionCwdConflict(sessionId, cwd, stored.cwd)
- }
- const inspected = await persistence.inspect(sessionId)
- if (hasSubagentOwner({ header: inspected.meta, events: inspected.events }, undefined)) {
- throw new SubagentSessionOwnership(sessionId)
- }
- return (await ctx.agents.resume({
- resumeSessionId: sessionId,
- agentOptions,
- setup: installTarget,
- })).agent
- }
- try {
- await mkdir(cwd, { recursive: true })
- } catch (error: unknown) {
- throw new Error(`failed to ensure project directory "${cwd}": ${String(error)}`, { cause: error })
- }
- return (await ctx.agents.create({
- sessionId,
- agentOptions,
- meta: { cwd },
- setup: installTarget,
- })).agent
- })().catch((error: unknown) => {
- // Another Host entry path may have published the same identity while
- // this operation crossed an asynchronous persistence/filesystem step.
- const live = ctx.agents.get(sessionId)
- if (live !== undefined) {
- if (hasSubagentOwner(live.session, live)) throw new SubagentSessionOwnership(sessionId)
- return live
- }
- const attached = ctx.sessions.get(sessionId)
- if (attached !== undefined && hasSubagentOwner(attached, undefined)) {
- throw new SubagentSessionOwnership(sessionId)
- }
- throw error
- }).finally(() => {
- sessionCreations.delete(sessionId)
- })
- sessionCreations.set(sessionId, creation)
- }
- const agent = await creation
- if (hasSubagentOwner(agent.session, agent)) throw new SubagentSessionOwnership(sessionId)
- if (agent.session.header.cwd !== cwd) {
- throw new SessionCwdConflict(sessionId, cwd, agent.session.header.cwd)
- }
- return agent
- }
- /** Resolve or create one path while holding the Host's workspace-create chain. */
- function ensureWorkspace(
- path: string,
- title: string | undefined,
- rejectExistingName = false,
- createDirectory = false,
- ): Promise<{ workspace: Workspace; created: boolean }> {
- const operation = workspaceCreationChain.then(async () => {
- if (rejectExistingName && title !== undefined
- && ctx.workspace.list().some(workspace => workspace.title === title)) {
- throw new WorkspaceNameConflictError(title)
- }
- if (createDirectory) {
- try {
- await mkdir(path, { recursive: true })
- } catch (error: unknown) {
- throw new WorkspaceDirectoryCreationError(
- `failed to create workspace directory "${path}": ${String(error)}`,
- )
- }
- }
- const existing = await ctx.workspace.resolveByPath(path)
- if (existing !== undefined) return { workspace: existing, created: false }
- return { workspace: await ctx.workspace.create(path, title), created: true }
- })
- workspaceCreationChain = operation.then(() => undefined, () => undefined)
- return operation
- }
- /**
- * Build the session.list baseline shared by listing and search visibility.
- * Attached sessions come from memory; servable cold sessions merge from
- * persistence, and the final order is newest-first.
- */
- async function listVisibleSessionSummaries(signal?: AbortSignal): Promise<SessionSummary[]> {
- signal?.throwIfAborted()
- const items = ctx.sessions.list().map((session) => {
- const agent = ctx.agents.get(session.id)
- const projections = listProjectionsFor(ctx, session.header, session)
- return {
- ...summarize(session, agent?.status === 'running'),
- ...projections === undefined ? {} : { projections },
- }
- })
- signal?.throwIfAborted()
- const attached = new Set(items.map(item => item.sessionId))
- const persistence = ctx.get('sessionPersistence')
- if (persistence !== undefined) {
- const cold = (await persistence.list(signal))
- .filter(meta => !attached.has(meta.id) && meta.cwd !== undefined)
- signal?.throwIfAborted()
- for (let offset = 0; offset < cold.length; offset += COLD_SUMMARY_BATCH_SIZE) {
- signal?.throwIfAborted()
- const batch = cold.slice(offset, offset + COLD_SUMMARY_BATCH_SIZE)
- const settled = await Promise.allSettled(
- batch.map(async (meta) => {
- // Cold rows read the persisted projection cache only — never a
- // log load; a session without a cache row simply has no column.
- const projections = listProjectionsFor(ctx, meta, undefined)
- return {
- ...await summarizeCold(persistence, meta, signal),
- ...projections === undefined ? {} : { projections },
- }
- }),
- )
- const summaries: SessionSummary[] = []
- let rejected = false
- let failure: unknown
- for (const result of settled) {
- if (result.status === 'fulfilled') {
- summaries.push(result.value)
- } else if (!rejected) {
- rejected = true
- failure = result.reason
- }
- }
- if (rejected) throw failure
- signal?.throwIfAborted()
- items.push(...summaries)
- }
- }
- items.sort((a, b) => b.updatedAt - a.updatedAt)
- return items
- }
- /** Resolve the goal service; absent = the deployment did not compose @deepseek-ai/dsh-goal. */
- function goalService(): NonNullable<ReturnType<typeof ctx.get<'goals'>>> | { error: RpcError } {
- const goals = ctx.get('goals')
- if (goals === undefined) {
- 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: {} } }
- }
- return goals
- }
- /** Map one goal-domain rejection to the wire error (stable GoalError codes ride in details). */
- function goalError(request: RpcRequest<unknown>, error: unknown): RpcResponse<never> {
- const details = error instanceof GoalError ? { goalCode: error.code } : {}
- return err(request, { code: 'internal', message: String(error), details })
- }
- /** Resolve a session's agent, apply one goal mutation, and acknowledge with the new CAS ref. */
- async function mutateGoal(
- request: RpcRequest<{ sessionId: SessionId }>,
- mutation: (goals: NonNullable<ReturnType<typeof ctx.get<'goals'>>>, agent: Agent) => CoreGoalRef,
- ): Promise<RpcResponse<{ ref: GoalRef }>> {
- const goals = goalService()
- if ('error' in goals) return err(request, goals.error)
- const found = await agentFor(request.payload.sessionId)
- if ('error' in found) return err(request, found.error)
- try {
- const ref = mutation(goals, found.agent)
- return ok(request, { ref: { id: ref.id, revision: ref.revision } })
- } catch (error: unknown) {
- return goalError(request, error)
- }
- }
- /** Missing-service report shared by the settings domain (skills-domain stance). */
- function settingsAbsent(): RpcError {
- 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: {} }
- }
- /** Missing-service report shared by the credentials domain. */
- function credentialsAbsent(): RpcError {
- 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: {} }
- }
- /** Map one redacted seam descriptor to its wire view. */
- function namespaceView(descriptor: SettingsDescriptor): SettingsNamespaceView {
- return {
- ns: String(descriptor.ns),
- schema: descriptor.schema,
- value: descriptor.value,
- ...descriptor.base === undefined ? {} : { base: descriptor.base },
- ...descriptor.user === undefined ? {} : { user: descriptor.user },
- applies: descriptor.applies,
- secrets: (descriptor.secrets ?? []).map(secret => ({ path: [...secret.path], set: secret.set })),
- revision: descriptor.revision,
- }
- }
- /** Settings namespaces whose changes can invalidate the model catalog. */
- function modelProviderNamespaces(): Set<string> {
- return new Set(ctx.llm.listConfigurableProviders().map(entry => entry.settingsNs))
- }
- /**
- * The settings namespaces this proxy serves: configurable model providers
- * plus the small explicit Web preference and product-owned allowlists. The
- * settings seam remains general; a future registration does not become
- * remotely readable or writable by default.
- */
- function exposedNamespaces(): Set<string> {
- const exposed = modelProviderNamespaces()
- for (const ns of WEB_SETTINGS_NAMESPACES) exposed.add(ns)
- for (const ns of PRODUCT_SETTINGS_NAMESPACES) exposed.add(ns)
- return exposed
- }
- /** Refuse a namespace outside the explicit configuration-client boundary. */
- function notExposed(request: RpcRequest<unknown>, ns: string): RpcResponse<SettingsNamespaceView> {
- return err(request, {
- code: 'settings-not-exposed',
- message: `settings namespace "${ns}" is not exposed to configuration clients`,
- details: { ns },
- })
- }
- /**
- * Run one settings write (merge or wholesale replace) and acknowledge with
- * the namespace's new redacted view. A namespace outside the configuration
- * boundary is refused before the seam is touched; every seam refusal —
- * unknown or invalid namespace, read-only provider, schema validation,
- * storage — becomes one `settings-rejected` carrying the seam's own message.
- */
- async function settingsWrite(
- request: RpcRequest<unknown>,
- ns: string,
- mode: 'update' | 'replace' | 'mutate',
- section: object,
- expectedRevision?: number,
- ): Promise<RpcResponse<SettingsNamespaceView>> {
- const settings = ctx.get('settings')
- if (settings === undefined) return err(request, settingsAbsent())
- const rejected = (error: unknown): RpcResponse<SettingsNamespaceView> => {
- // A stale writer is its own outcome, not a malformed request: the client
- // must re-read and re-apply rather than treat the write as invalid.
- if (error instanceof SettingsConflictError) {
- return err(request, {
- code: 'settings-conflict',
- message: error.message,
- details: { ns, expected: error.expected, actual: error.actual },
- })
- }
- return err(request, {
- code: 'settings-rejected',
- message: error instanceof Error ? error.message : String(error),
- details: { ns },
- })
- }
- let branded: SettingsNamespace
- try {
- branded = settingsNamespace(ns)
- } catch (error: unknown) {
- // A malformed name is a client bug, reported as such; it could never be
- // in the exposed set either, so naming the real fault costs no ground.
- return rejected(error)
- }
- if (!exposedNamespaces().has(ns)) return notExposed(request, ns)
- try {
- if (mode === 'update') await settings.update(branded, section, expectedRevision)
- else if (mode === 'replace') await settings.replace(branded, section, expectedRevision)
- else await settings.mutate(branded, section as SettingsPathOp[], expectedRevision)
- } catch (error: unknown) {
- return rejected(error)
- }
- const descriptor = settings.describe({ redactSecrets: true }).find(candidate => candidate.ns === branded)
- if (descriptor === undefined) {
- // The write committed but the namespace vanished before this read: only
- // a concurrent registrant disposal can produce it.
- return err(request, { code: 'internal', message: `settings namespace "${ns}" was disposed after the ${mode}`, details: {} })
- }
- return ok(request, namespaceView(descriptor))
- }
- return {
- sessions: {
- // Attached sessions summarize from memory; persisted-but-unattached (cold)
- // sessions merge in from the persistence store so history survives restarts.
- // Legacy logs without a cwd (pre-project stance) are not served — every
- // session now records its project at create time.
- async list(request) {
- return ok(request, { items: await listVisibleSessionSummaries() })
- },
- async search(request, signal) {
- const cancelled = () => err<{ items: SessionSearchItem[]; hasMore: boolean }>(request, {
- code: 'cancelled',
- message: 'session search was aborted',
- details: {},
- })
- if (isAborted(signal)) return cancelled()
- const sessionQuery = ctx.get('sessionQuery')
- if (sessionQuery === undefined) {
- return err(request, {
- code: 'internal',
- message: 'session search is unavailable: this deployment does not mount @deepseek-ai/dsh-session-query',
- details: {},
- })
- }
- try {
- const visible = await listVisibleSessionSummaries(signal)
- if (isAborted(signal)) return cancelled()
- if (visible.length === 0) return ok(request, { items: [], hasMore: false })
- const visibleIds = new Set(visible.map(item => item.sessionId))
- const authorized: SessionSearchItem[] = []
- const acceptedIds = new Set<SessionId>()
- const seenCursors = new Set<SessionSearchCursor>()
- let cursor: SessionSearchCursor | undefined
- let providerCallCount = 0
- let providerPageLimit = SESSION_SEARCH_RESULT_LIMIT
- while (authorized.length <= SESSION_SEARCH_RESULT_LIMIT) {
- if (isAborted(signal)) return cancelled()
- if (providerCallCount >= SESSION_SEARCH_PROVIDER_CALL_LIMIT) {
- throw new Error(
- `session search provider exceeded the ${SESSION_SEARCH_PROVIDER_CALL_LIMIT}-call work budget`,
- )
- }
- providerCallCount++
- const requestedCursor = cursor
- const requestedPageLimit = providerPageLimit
- let page
- try {
- page = await sessionQuery.searchSessions({
- query: request.payload.query,
- eventFilters: [
- { kind: 'type', values: ['user/message', 'assistant/message', 'steering/message'] },
- { kind: 'surface', values: ['current'] },
- ],
- limit: requestedPageLimit,
- ...requestedCursor === undefined ? {} : { cursor: requestedCursor },
- }, { signal })
- } catch (error: unknown) {
- if (isAborted(signal)) return cancelled()
- if (
- requestedCursor === undefined
- && error instanceof SessionQueryError
- && error.code === 'SESSION_QUERY_INVALID_LIMIT'
- && requestedPageLimit > 1
- ) {
- providerPageLimit = Math.max(1, Math.floor(requestedPageLimit / 2))
- continue
- }
- if (
- requestedCursor !== undefined
- && error instanceof SessionQueryError
- && error.code === 'SESSION_QUERY_STALE_CURSOR'
- ) {
- authorized.length = 0
- acceptedIds.clear()
- seenCursors.clear()
- cursor = undefined
- continue
- }
- throw error
- }
- if (isAborted(signal)) return cancelled()
- const providerItemCount = page.items.length
- if (providerItemCount > requestedPageLimit) {
- throw new Error(
- `session search provider returned ${providerItemCount} items; maximum is ${requestedPageLimit}`,
- )
- }
- // Host visibility is the authorization boundary. Consume the
- // provider's globally ranked stream rather than binding every
- // visible id into one SQLite statement, then re-check complete
- // provenance before emitting any snippet.
- for (const hit of page.items) {
- if (authorized.length > SESSION_SEARCH_RESULT_LIMIT) continue
- if (
- !visibleIds.has(hit.header.id)
- || hit.bestMatch.sessionId !== hit.header.id
- || hit.bestMatch.surface !== 'current'
- || !MESSAGE_TYPES.has(hit.bestMatch.type)
- || acceptedIds.has(hit.header.id)
- ) continue
- const snippet = truncateUnicodeCodePoints(
- hit.bestMatch.snippet,
- SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS,
- )
- acceptedIds.add(hit.header.id)
- authorized.push({
- sessionId: hit.header.id,
- snippet,
- })
- }
- const nextCursor = page.nextCursor
- if (nextCursor !== undefined) {
- if (seenCursors.has(nextCursor)) {
- throw new Error('session search provider repeated a continuation cursor')
- }
- seenCursors.add(nextCursor)
- }
- if (authorized.length > SESSION_SEARCH_RESULT_LIMIT || nextCursor === undefined) break
- cursor = nextCursor
- }
- return ok(request, {
- items: authorized.slice(0, SESSION_SEARCH_RESULT_LIMIT),
- hasMore: authorized.length > SESSION_SEARCH_RESULT_LIMIT,
- })
- } catch (error: unknown) {
- if (
- isAborted(signal)
- || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')
- ) return cancelled()
- // XXX: Redact provider details before exposing this gateway beyond
- // its current single-user local deployment.
- return err(request, {
- code: 'internal',
- message: `session search failed: ${String(error)}`,
- details: {},
- })
- }
- },
- async create(request) {
- const sessionId = request.payload.sessionId ?? `session-${randomUUID()}` as SessionId
- let workspace: Workspace | undefined
- if (request.payload.workspaceId !== undefined) {
- workspace = ctx.workspace.get(brandWorkspaceId(request.payload.workspaceId))
- if (workspace === undefined) {
- return err(request, {
- code: 'workspace-not-found',
- message: `workspace "${request.payload.workspaceId}" not found`,
- details: { workspaceId: request.payload.workspaceId },
- })
- }
- }
- const cwd = workspace?.path ?? request.payload.cwd ?? defaults.cwd
- try {
- await ensureSession(sessionId, cwd, request.payload.sessionId !== undefined)
- } catch (error: unknown) {
- if (error instanceof SessionCwdConflict) {
- return err(request, {
- code: 'session-conflict',
- message: error.message,
- details: {
- sessionId: error.sessionId,
- requestedCwd: error.requestedCwd,
- ...error.existingCwd === undefined ? {} : { existingCwd: error.existingCwd },
- },
- })
- }
- if (error instanceof SubagentSessionOwnership) {
- return err(request, subagentOwnershipError(error.sessionId))
- }
- return err(request, {
- code: 'internal',
- message: `failed to create session "${sessionId}": ${String(error)}`,
- details: {},
- })
- }
- if (workspace !== undefined) {
- try {
- await workspace.attachSession(sessionId)
- } catch (error: unknown) {
- return err(request, {
- code: 'workspace-attach-failed',
- message: `session "${sessionId}" was created but could not attach to workspace "${workspace.id}": ${String(error)}`,
- details: { sessionId, workspaceId: workspace.id },
- })
- }
- }
- return ok(request, { sessionId })
- },
- async history(request) {
- const { sessionId, beforeSeq, maxMessages } = request.payload
- let state: { events: SessionEvent[]; projections?: SessionProjectionsBlock }
- try {
- state = await historyStateFor(sessionId, beforeSeq === undefined)
- } catch (error: unknown) {
- if (error instanceof SessionNotFound) {
- return err(request, { code: 'session-not-found', message: error.message, details: { sessionId } })
- }
- return err(request, {
- code: 'internal',
- message: `history unavailable for session "${sessionId}": ${String(error)}`,
- details: {},
- })
- }
- const page = historyPage(ctx, state.events, beforeSeq, maxMessages)
- return ok(request, {
- events: page.events,
- hasMore: page.hasMore,
- ...state.projections === undefined ? {} : { projections: state.projections },
- })
- },
- async models(request) {
- const { sessionId } = request.payload
- const found = await agentFor(sessionId)
- if ('error' in found) return err(request, found.error)
- const current = targetFor(found.agent).current
- const { groups, failures } = await buildModelCatalog(ctx, current)
- return ok(request, { current: { ...current }, groups, failures })
- },
- async selectModel(request) {
- const { sessionId, provider, model, reasoningEffort } = request.payload
- const found = await agentFor(sessionId)
- if ('error' in found) return err(request, found.error)
- try {
- const resolved = await ctx.llm.resolveCallConfig({
- provider,
- model,
- ...reasoningEffort === undefined
- ? {}
- : { reasoningEffort: ReasoningEffortId(reasoningEffort) },
- })
- const selected: AgentLlmTarget = {
- provider: resolved.provider,
- model: resolved.model,
- ...resolved.reasoningEffort === undefined
- ? {}
- : { reasoningEffort: resolved.reasoningEffort },
- }
- targetFor(found.agent).current = selected
- return ok(request, { selected: { ...selected } })
- } catch (error: unknown) {
- return err(request, {
- code: 'model-unavailable',
- message: error instanceof Error ? error.message : String(error),
- details: { provider, model },
- })
- }
- },
- async rename(request) {
- const { sessionId, title } = request.payload
- const found = await agentFor(sessionId)
- if ('error' in found) return err(request, found.error)
- const titles = ctx.get('sessionTitle')
- if (titles === undefined) {
- return err(request, { code: 'internal', message: 'renaming is unavailable: this deployment mounts no session-title service', details: {} })
- }
- try {
- const accepted = titles.rename(found.agent.session, title)
- return ok(request, { title: accepted.title, seq: accepted.eventSeq })
- } catch (error: unknown) {
- // Only the input's fault maps to title-invalid (the message is
- // product-user-visible in the rename dialog); liveness and disposal
- // races are deployment trouble, not a bad title.
- if (error instanceof SessionTitleInvalidError) {
- return err(request, {
- code: 'title-invalid',
- message: error.message,
- details: { sessionId },
- })
- }
- return err(request, {
- code: 'internal',
- message: `failed to rename session "${sessionId}": ${String(error)}`,
- details: {},
- })
- }
- },
- async fork(request) {
- const { sessionId, atSeq } = request.payload
- let source: SessionReadState
- try {
- source = await readSessionState(sessionId)
- } catch (error: unknown) {
- if (error instanceof SessionNotFound) {
- return err(request, { code: 'session-not-found', message: error.message, details: { sessionId } })
- }
- return err(request, {
- code: 'internal',
- message: `fork source unavailable for session "${sessionId}": ${String(error)}`,
- details: {},
- })
- }
- const events = source.events
- // An in-log anchor belongs to the turn containing it and must never
- // clip backward to an earlier completed turn. Omitted and past-end
- // anchors retain the last-completed-turn shortcut.
- const lastSeq = events.at(-1)?.seq ?? -1
- const anchoredBoundary = atSeq === undefined
- ? undefined
- : events.find(e => e.type === 'turn/end' && e.seq >= atSeq)
- const boundary = anchoredBoundary
- ?? (atSeq === undefined || atSeq > lastSeq
- ? events.findLast(e => e.type === 'turn/end')
- : undefined)
- if (boundary === undefined) {
- return err(request, {
- code: 'fork-unavailable',
- message: atSeq !== undefined && atSeq <= lastSeq
- ? `session "${sessionId}" has not completed the turn containing event ${String(atSeq)}`
- : `session "${sessionId}" has no completed turn to fork from`,
- details: { sessionId },
- })
- }
- // Extend the cut through trailing out-of-band appends (session/title,
- // injections) up to the next turn/start: they are standalone events, so
- // the seed stays balanced, and the child inherits a title generated
- // right after the boundary turn.
- let cut = boundary.seq + 1
- while (cut < events.length && events[cut]?.type !== 'turn/start') cut++
- let workspace: Workspace | undefined
- try {
- workspace = await forkWorkspace(source)
- } catch (error: unknown) {
- return err(request, {
- code: 'internal',
- message: `failed to resolve fork workspace for session "${sessionId}": ${String(error)}`,
- details: {},
- })
- }
- const childId = `session-${randomUUID()}` as SessionId
- try {
- await ctx.agents.create({
- sessionId: childId,
- seed: events.slice(0, cut),
- meta: {
- ...source.header.cwd === undefined ? {} : { cwd: source.header.cwd },
- parentSession: source.id,
- seedLength: cut,
- },
- agentOptions,
- setup: installTarget,
- })
- } catch (error: unknown) {
- return err(request, {
- code: 'internal',
- message: `failed to fork session "${sessionId}": ${String(error)}`,
- details: {},
- })
- }
- // An ordinary source keeps its direct Workspace. A subagent source is
- // not listed there, so its ordinary fork joins the nearest owning
- // ancestor instead. The child is already published if attach fails.
- if (workspace !== undefined) {
- try {
- await workspace.attachSession(childId)
- } catch (error: unknown) {
- return err(request, {
- code: 'workspace-attach-failed',
- message: `session "${childId}" was forked but could not attach to workspace "${workspace.id}": ${String(error)}`,
- details: { sessionId: childId, workspaceId: workspace.id },
- })
- }
- }
- return ok(request, { sessionId: childId })
- },
- async prompt(request) {
- const { sessionId, mode, content } = request.payload
- const found = await agentFor(sessionId)
- if ('error' in found) return err(request, found.error)
- const agent = found.agent
- // The rpcId rides MessageSource into user/message (merge declaration in api/sessions.ts; provisional correlation).
- const source: MessageSource = { kind: 'user', rpcId: request.rpcId }
- try {
- const message: UserMessage = createUserMessage({ content, source })
- if (mode === 'steer') agent.steer(message)
- else agent.followup(message)
- } catch (error: unknown) {
- // A synchronous throw from steer/followup means disposed or invalid input; surface as agent-busy with the reason attached.
- return err(request, { code: 'agent-busy', message: 'prompt rejected', details: { reason: String(error) } })
- }
- return ok(request, { accepted: true as const })
- },
- updateQueue(request) {
- const { sessionId, itemId, action } = request.payload
- const agent = ctx.agents.get(sessionId)
- if (agent !== undefined && hasSubagentOwner(agent.session, agent)) {
- return Promise.resolve(err(request, subagentOwnershipError(sessionId)))
- }
- if (agent === undefined || agent.updateInbox(itemId, action) === 'not-found') {
- return Promise.resolve(err(request, {
- code: 'queue-item-not-found',
- message: 'queued item is no longer pending',
- details: { itemId },
- }))
- }
- return Promise.resolve(ok(request, { accepted: true as const }))
- },
- cancel(request) {
- const { sessionId } = request.payload
- const agent = ctx.agents.get(sessionId)
- if (agent === undefined) {
- return Promise.resolve(err(request, {
- code: 'session-not-found',
- message: `session "${sessionId}" not found (not attached)`,
- details: { sessionId },
- }))
- }
- if (hasSubagentOwner(agent.session, agent)) {
- return Promise.resolve(err(request, subagentOwnershipError(sessionId)))
- }
- agent.cancel({ kind: 'user' }, { keepInbox: true })
- return Promise.resolve(ok(request, { accepted: true as const }))
- },
- },
- subagents: {
- async list(request, signal) {
- try {
- const entries = await ctx.subagents.listChildren(request.payload.parentSessionId, signal)
- return ok(request, {
- entries: entries.map(entry => entry.kind === 'child'
- ? {
- ...entry,
- activity: ctx.agents.get(entry.id)?.status === 'running' ? 'running' : 'inactive',
- }
- : entry),
- parentAvailable: ctx.agents.get(request.payload.parentSessionId) !== undefined,
- })
- } catch (error: unknown) {
- if (signal?.aborted
- || (error instanceof SubagentError && error.code === 'CANCELLED')
- || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')) {
- return err(request, {
- code: 'cancelled',
- message: 'subagent catalog read was cancelled',
- details: {},
- })
- }
- return err(request, {
- code: 'internal',
- message: 'subagent catalog read failed',
- details: {},
- })
- }
- },
- async history(request, signal) {
- const {
- parentSessionId, childSessionId, mode, beforeSeq, maxMessages,
- } = request.payload
- const verified = await catalogChild(ctx, {
- parentSessionId, childSessionId, mode,
- }, signal)
- if (verified.error !== undefined) return err(request, verified.error)
- try {
- const snapshot = await ctx.sessionQuery.readSession(childSessionId)
- signal?.throwIfAborted()
- if (snapshot.session.parentSession !== parentSessionId) {
- return err(request, {
- code: 'subagent-unauthorized',
- message: 'subagent parent changed during history read',
- details: { childSessionId },
- })
- }
- const page = historyPage(ctx, snapshot.events, beforeSeq, maxMessages)
- const projections = beforeSeq === undefined
- ? detachedProjectionsFor(ctx, snapshot.events)
- : undefined
- return ok(request, { ...page, ...projections === undefined ? {} : { projections } })
- } catch (error: unknown) {
- if (signal?.aborted
- || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')) {
- return err(request, {
- code: 'cancelled',
- message: 'subagent history read was cancelled',
- details: {},
- })
- }
- if (error instanceof SessionQueryError
- && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') {
- return err(request, {
- code: 'subagent-not-found',
- message: 'subagent disappeared during history read',
- details: { parentSessionId, childSessionId },
- })
- }
- return err(request, {
- code: 'internal',
- message: 'subagent history read failed',
- details: {},
- })
- }
- },
- async prompt(request, signal) {
- const { parentSessionId, childSessionId, content } = request.payload
- const parent = ctx.agents.get(parentSessionId)
- if (parent === undefined) {
- return err(request, {
- code: 'subagent-parent-unavailable',
- message: `parent session "${parentSessionId}" is not live`,
- details: { parentSessionId },
- })
- }
- const verified = await catalogChild(ctx, {
- parentSessionId, childSessionId, mode: 'continuable',
- }, signal)
- if (verified.error !== undefined) return err(request, verified.error)
- try {
- const messageId = await ctx.subagents.followup(parent, childSessionId, content, {
- source: { kind: 'user', rpcId: request.rpcId },
- signal,
- })
- return ok(request, { messageId })
- } catch (error: unknown) {
- return subagentPromptError(request, error, signal)
- }
- },
- },
- workspace: {
- list(request) {
- return Promise.resolve(ok(request, {
- items: ctx.workspace.list().map(workspaceView),
- archivedSessionIds: [...ctx.workspace.archivedSessionIds],
- }))
- },
- // Exactly one of path/name arrives (schema refine). Existing-folder
- // adoption reuses its canonical path; create-by-name rejects a name
- // already present in the registry.
- // TODO: the create-by-name branch lost its last product consumer when
- // the Web picker collapsed onto the directory flow
- // (.agents/notes/implemented/simplification/2026-07-31-one-route-to-add-a-workspace.md).
- // Delete it with the wire schema's `name` member, this
- // `defaults.workspaceRoot`, the client seam that carried the name
- // (`WorkspaceCreateInput`, `WorkspacesService.create`'s `{ name }` arm,
- // `intentName`'s name branch, the manager's "name under workspaceRoot"
- // contract), and the `dsh web --workspace-root` flag plus its apps/cli
- // README lines, which exist only to feed it.
- async create(request) {
- const { payload } = request
- let path: string
- if (payload.name !== undefined) {
- const name = payload.name.trim()
- if (name === '' || name === '.' || name === '..' || /[/\\]/.test(name)) {
- return err(request, {
- code: 'workspace-invalid-path',
- message: `workspace name must be one non-empty path segment, got "${payload.name}"`,
- details: { path: payload.name },
- })
- }
- path = join(defaults.workspaceRoot, name)
- } else {
- path = payload.path as string
- }
- try {
- const name = payload.name?.trim()
- const { workspace, created } = await ensureWorkspace(
- path,
- name,
- name !== undefined,
- name !== undefined,
- )
- return ok(request, { workspace: workspaceView(workspace), created })
- } catch (error: unknown) {
- if (error instanceof WorkspaceNameConflictError) {
- return err(request, {
- code: 'workspace-name-conflict',
- message: error.message,
- details: { name: error.workspaceName },
- })
- }
- if (error instanceof WorkspaceDirectoryCreationError) {
- return err(request, { code: 'internal', message: error.message, details: {} })
- }
- // The registry rejects a path that does not resolve to an existing
- // directory (realpath ENOENT / not-a-directory) — the business
- // error of the typed-path flow, surfaced as a validation failure.
- return err(request, {
- code: 'workspace-invalid-path',
- message: `cannot create a workspace at "${path}": ${error instanceof Error ? error.message : String(error)}`,
- details: { path },
- })
- }
- },
- async rename(request) {
- const { payload } = request
- const workspace = ctx.workspace.get(brandWorkspaceId(payload.workspaceId))
- if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
- const title = payload.title.trim()
- // Uniqueness AND the same-title no-op both ride the create chain so
- // they observe the state left by earlier queued renames — checked
- // up front, a queued A→A could report success while an earlier A→B
- // still lands afterwards.
- const operation = workspaceCreationChain.then(async () => {
- if (title === workspace.title) return
- if (ctx.workspace.list().some(other => other.id !== workspace.id && other.title === title)) {
- throw new WorkspaceNameConflictError(title)
- }
- await workspace.setTitle(title)
- })
- workspaceCreationChain = operation.then(() => undefined, () => undefined)
- try {
- await operation
- } catch (error: unknown) {
- if (error instanceof WorkspaceNameConflictError) {
- return err(request, {
- code: 'workspace-name-conflict',
- message: error.message,
- details: { name: error.workspaceName },
- })
- }
- throw error
- }
- return ok(request, { workspace: workspaceView(workspace) })
- },
- async delete(request) {
- const { workspaceId } = request.payload
- const operation = workspaceCreationChain.then(() =>
- ctx.workspace.delete(brandWorkspaceId(workspaceId)))
- workspaceCreationChain = operation.then(() => undefined, () => undefined)
- if (!await operation) return workspaceNotFound(request, workspaceId)
- return ok(request, { deleted: true as const })
- },
- async insertSessionBefore(request) {
- const { payload } = request
- const workspace = ctx.workspace.get(brandWorkspaceId(payload.workspaceId))
- if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
- try {
- await workspace.insertSessionBefore(payload.sessionId, payload.beforeSessionId)
- } catch (error: unknown) {
- // Only the entity's unaccounted-id rejection is the business code;
- // storage/durability failures propagate as internal errors.
- if (!(error instanceof WorkspaceMoveInvalidError)) throw error
- return err(request, {
- code: 'workspace-move-invalid',
- message: error.message,
- details: {
- workspaceId: payload.workspaceId,
- sessionId: payload.sessionId,
- ...payload.beforeSessionId === undefined ? {} : { beforeSessionId: payload.beforeSessionId },
- },
- })
- }
- return ok(request, { workspace: workspaceView(workspace) })
- },
- async archiveSession(request) {
- const { sessionId } = request.payload
- try {
- await ctx.workspace.archiveSession(sessionId)
- } catch (error: unknown) {
- // Only the registry's unknown-session rejection is the business
- // code; storage/durability failures propagate as internal errors.
- if (!(error instanceof WorkspaceUnknownSessionError)) throw error
- return err(request, {
- code: 'session-not-found',
- message: error.message,
- details: { sessionId },
- })
- }
- return ok(request, { archivedSessionIds: [...ctx.workspace.archivedSessionIds] })
- },
- },
- host: {
- describe(request) {
- // TODO(step2): version should read apps/cli's package.json; placeholder for now.
- return Promise.resolve(ok(request, {
- version: '0.0.1',
- // Same source as session.create's fallback: the UI's default project
- // must match where an unspecified-cwd session actually lands.
- cwd: defaults.cwd,
- provider: defaults.provider,
- model: defaults.model,
- attachedSessions: ctx.agents.list().length,
- }))
- },
- async pickDirectory(request, signal) {
- const capability = ctx.directoryPicker.capability()
- if (capability.kind !== 'native') {
- return err(request, {
- code: 'directory-picker-unavailable',
- message: `host.pickDirectory needs the native capability; the composed picker serves "${capability.kind}"`,
- details: { capability: capability.kind },
- })
- }
- try {
- const path = await capability.pick(signal)
- return ok(request, { path })
- } catch (error: unknown) {
- if (signal.aborted) {
- return err(request, {
- code: 'cancelled',
- message: 'directory picker was aborted',
- details: {},
- })
- }
- return err(request, {
- code: 'internal',
- message: `directory picker failed: ${error instanceof Error ? error.message : String(error)}`,
- details: {},
- })
- }
- },
- async listDirectory(request, signal) {
- const capability = ctx.directoryPicker.capability()
- if (capability.kind !== 'browse') {
- return err(request, {
- code: 'directory-picker-unavailable',
- message: `host.listDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
- details: { capability: capability.kind },
- })
- }
- try {
- // The carrier's signal follows the caller: a disconnect or timeout
- // stops the backend's directory scan instead of outliving it.
- return ok(request, await capability.list(request.payload.path, signal))
- } catch (error: unknown) {
- // An abort is the caller's own timeout/disconnect, not a server
- // failure — same code pickDirectory and command.execute report.
- if (signal.aborted) {
- return err(request, { code: 'cancelled', message: 'directory listing was aborted', details: {} })
- }
- return err(request, directoryError(error))
- }
- },
- async createDirectory(request) {
- const capability = ctx.directoryPicker.capability()
- if (capability.kind !== 'browse') {
- return err(request, {
- code: 'directory-picker-unavailable',
- message: `host.createDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
- details: { capability: capability.kind },
- })
- }
- try {
- return ok(request, { path: await capability.createDirectory(request.payload.path, request.payload.name) })
- } catch (error: unknown) {
- return err(request, directoryError(error))
- }
- },
- async openPath(request, signal) {
- try {
- const open = defaults.openPath
- ?? ((path: string, openSignal: AbortSignal) => openNativePath(path, openSignal))
- await open(request.payload.path, signal)
- return ok(request, { opened: true as const })
- } catch (error: unknown) {
- if (signal.aborted) {
- return err(request, {
- code: 'cancelled',
- message: 'path open was aborted',
- details: {},
- })
- }
- return err(request, {
- code: 'internal',
- message: `path open failed: ${error instanceof Error ? error.message : String(error)}`,
- details: {},
- })
- }
- },
- },
- commands: {
- // Both methods address one session's agent (agentFor keeps its
- // resume-on-miss: clients only send a sessionId for a published
- // session, and resume restores an existing entity).
- async list(request) {
- // Missing service = the deployment omitted dsh-commands from its
- // composition, not an empty catalog: fail loud instead of serving [].
- const commands = ctx.get('commands')
- if (commands === undefined) {
- 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: {} })
- }
- const found = await agentFor(request.payload.sessionId)
- if ('error' in found) return err(request, found.error)
- return ok(request, { commands: commands.list(found.agent) })
- },
- async execute(request, signal) {
- const commands = ctx.get('commands')
- if (commands === undefined) {
- 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: {} })
- }
- const { sessionId, line } = request.payload
- const found = await agentFor(sessionId)
- if ('error' in found) return err(request, found.error)
- try {
- // Pure admission: the executor's durable command/run + command/done
- // pair (broadcast on the mux stream) carries the outcome; the
- // response reports whether the line resolved to a handler, plus the
- // minted pairing id so the issuing client can correlate its request
- // with the flow node the lifecycle events produce.
- const execution = await commands.execute(found.agent, line, signal)
- return ok(request, execution === undefined
- ? { matched: false }
- : { matched: true, commandId: execution.commandId })
- } catch (error: unknown) {
- if (signal.aborted) return err(request, { code: 'cancelled', message: 'command execution was aborted', details: {} })
- return err(request, { code: 'internal', message: `command failed: ${String(error)}`, details: {} })
- }
- },
- },
- goals: {
- // Mutations only — the read side is the 'goal' session projection.
- // Every verb resolves the session's agent (agentFor: implicit cold
- // resume, the command.* precedent) and acknowledges with the new CAS
- // ref; the committed goal/change event carries the whole value to every
- // client through the projection frames.
- async create(request) {
- const { objective, maxGoalRounds } = request.payload
- return mutateGoal(request, (goals, agent) => goals.create(agent, {
- objective,
- ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
- }))
- },
- async edit(request) {
- const { ref, objective, maxGoalRounds } = request.payload
- return mutateGoal(request, (goals, agent) => goals.edit(agent, ref, {
- ...(objective !== undefined ? { objective } : {}),
- ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
- }))
- },
- async pause(request) {
- return mutateGoal(request, (goals, agent) => goals.pause(agent, request.payload.ref))
- },
- async resume(request) {
- return mutateGoal(request, (goals, agent) => goals.resume(agent, request.payload.ref))
- },
- async complete(request) {
- return mutateGoal(request, (goals, agent) => goals.complete(agent, request.payload.ref))
- },
- async clear(request) {
- const goals = goalService()
- if ('error' in goals) return err(request, goals.error)
- const found = await agentFor(request.payload.sessionId)
- if ('error' in found) return err(request, found.error)
- try {
- goals.clear(found.agent, request.payload.ref)
- return ok(request, { cleared: true as const })
- } catch (error: unknown) {
- return goalError(request, error)
- }
- },
- },
- skills: {
- // Skill lookup never touches the Agent registry: the session address
- // resolves to a canonical cwd from the host-resident session header, so
- // listing skills cannot create or resume an agent as a side effect.
- async list(request) {
- const { sessionId } = request.payload
- const session = ctx.sessions.get(sessionId)
- if (session === undefined) {
- return err(request, {
- code: 'session-not-found',
- message: `session "${sessionId}" not found (not attached)`,
- details: { sessionId },
- })
- }
- if (session.header.cwd === undefined) {
- // Every served session records its project at create time; a
- // cwd-less header is a pre-project legacy log (not served).
- return err(request, { code: 'internal', message: `session "${sessionId}" has no project cwd`, details: {} })
- }
- const cwd = session.header.cwd
- // Same stance as the commands domain: a missing service means the
- // deployment omitted dsh-skill from its composition, not an empty
- // catalog. ctx.get also keeps this handler independent of the gateway
- // plugin's inject list (an undeclared `ctx.skills` property read
- // fails the reflect proxy).
- const skillRegistry = ctx.get('skills')
- if (skillRegistry === undefined) {
- 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: {} })
- }
- try {
- const skills = (await skillRegistry.list({ cwd }))
- .filter(skill => skill.invocation.modelInvocable && skill.invocation.userInvocable)
- return ok(request, {
- skills: skills.map(skill => ({
- name: skill.name,
- description: skill.description,
- ...skill.whenToUse === undefined ? {} : { whenToUse: skill.whenToUse },
- })),
- })
- } catch (error: unknown) {
- return err(request, { code: 'internal', message: `skill listing failed: ${String(error)}`, details: {} })
- }
- },
- },
- settings: {
- describe(request) {
- const settings = ctx.get('settings')
- if (settings === undefined) return Promise.resolve(err(request, settingsAbsent()))
- const exposed = exposedNamespaces()
- return Promise.resolve(ok(request, {
- writable: settings.writable,
- namespaces: settings.describe({ redactSecrets: true })
- .filter(descriptor => exposed.has(String(descriptor.ns)))
- .map(namespaceView),
- }))
- },
- update: request => settingsWrite(request, request.payload.ns, 'update', request.payload.patch, request.payload.expectedRevision),
- replace: request => settingsWrite(request, request.payload.ns, 'replace', request.payload.section, request.payload.expectedRevision),
- mutate: request => settingsWrite(request, request.payload.ns, 'mutate', request.payload.ops, request.payload.expectedRevision),
- },
- credentials: {
- async describe(request) {
- const credentials = ctx.get('credentials')
- if (credentials === undefined) return err(request, credentialsAbsent())
- const entries = await Promise.all(request.payload.refs.map(async (ref) => {
- const info = await credentials.describe(credentialRef(ref))
- const view: CredentialView = {
- configured: info.configured,
- ...info.source === undefined ? {} : { source: info.source },
- writable: info.writable,
- }
- return [ref, view] as const
- }))
- return ok(request, { credentials: Object.fromEntries(entries) })
- },
- async set(request) {
- const credentials = ctx.get('credentials')
- if (credentials === undefined) return err(request, credentialsAbsent())
- const { ref, value } = request.payload
- try {
- await credentials.set(credentialRef(ref), value)
- } catch (error: unknown) {
- return err(request, {
- code: 'credential-rejected',
- message: error instanceof Error ? error.message : String(error),
- details: { ref },
- })
- }
- return ok(request, {})
- },
- async unset(request) {
- const credentials = ctx.get('credentials')
- if (credentials === undefined) return err(request, credentialsAbsent())
- const { ref } = request.payload
- try {
- await credentials.unset(credentialRef(ref))
- } catch (error: unknown) {
- return err(request, {
- code: 'credential-rejected',
- message: error instanceof Error ? error.message : String(error),
- details: { ref },
- })
- }
- return ok(request, {})
- },
- },
- llm: {
- providers(request) {
- const registered = ctx.llm.listProviders()
- const active = new Set(registered.map(provider => provider.id))
- const directory = ctx.llm.listConfigurableProviders()
- const declared = new Set(directory.map(entry => entry.provider))
- const views = directory.map(entry => ({
- provider: entry.provider,
- displayName: entry.displayName,
- settingsNs: entry.settingsNs,
- settingsPath: [...entry.settingsPath],
- active: active.has(entry.provider),
- }))
- // Routes registered without a directory declaration still appear —
- // they exist and serve models — just with no settings address.
- for (const provider of registered) {
- if (declared.has(provider.id)) continue
- views.push({
- provider: provider.id,
- displayName: provider.name,
- settingsNs: '',
- settingsPath: [],
- active: true,
- })
- }
- return Promise.resolve(ok(request, { providers: views }))
- },
- async models(request) {
- return ok(request, await buildModelCatalog(ctx))
- },
- },
- events: {
- mux(_request, signal) {
- const queue = new FrameQueue<RpcRequest<MuxFrame>>()
- muxQueues.add(queue)
- for (const session of ctx.sessions.list()) {
- subscribeSession(queue, session)
- }
- for (const pending of pendingQuestions.values()) {
- queue.push({
- rpcId: pending.rpcId,
- payload: {
- type: 'question/requested', sessionId: pending.sessionId,
- questions: pending.questions,
- },
- })
- }
- // Refresh recovery: still-pending approval questions replay with their
- // stable rpcId so a reconnecting client can still answer them.
- for (const pending of pendingApprovals.values()) queue.push(requestedFrame(pending))
- // Queue snapshot baseline (pendingQuestions precedent): frames replayed
- // in arrival order per session; a reconnecting client rebuilds its
- // queue view from these alone.
- for (const [sessionId, items] of queuedMirror) {
- queue.push(frame({
- type: 'session/queue',
- sessionId,
- items: items.map(item => ({
- id: item.id,
- message: item.message,
- })),
- }))
- }
- // Per-session open-call table for result-view pairing. Bounded by the
- // per-turn call count: entries clear on turn/end; a table miss (stream
- // opened mid-turn) backscans the session's in-memory events instead.
- const openCalls = new Map<SessionId, Map<string, { name: string; args: unknown }>>()
- const disposers = [
- ctx.on('session/event', (session: Session, event: SessionEvent) => {
- if (event.type === 'tool/call') {
- const data = event.data as ToolCallData
- try {
- let table = openCalls.get(session.id)
- if (table === undefined) openCalls.set(session.id, table = new Map<string, { name: string; args: unknown }>())
- table.set(data.callId, { name: data.name, args: JSON.parse(data.arguments) })
- } catch {
- // Unparseable model arguments: leave the table unset; the result view soft-falls.
- }
- } else if (event.type === 'turn/end') {
- openCalls.delete(session.id)
- }
- const view = viewFor(ctx, event, callId =>
- openCalls.get(session.id)?.get(callId) ?? backscanArgs(session.events, callId))
- queue.push(frame({ type: 'session/event', sessionId: session.id, event, ...view === undefined ? {} : { view } }))
- }),
- ctx.on('session/created', (session: Session) => {
- subscribeSession(queue, session)
- }),
- ctx.on('session/disposed', (session: Session) => {
- openCalls.delete(session.id)
- }),
- ]
- return queue.iterate(signal, () => {
- muxQueues.delete(queue)
- for (const dispose of disposers) dispose()
- })
- },
- host(_request, signal) {
- const queue = new FrameQueue<RpcRequest<HostFrame>>()
- const committedWorkspaceIds = new Set(
- ctx.workspace.list().map(workspace => String(workspace.id)),
- )
- // Frame-dedup baseline, same posture as committedWorkspaceIds: the
- // stream opens against the current set; workspace.list re-baselines
- // reconnecting clients, so only later changes need frames.
- let archivedSessionIds = ctx.workspace.archivedSessionIds
- const disposers = [
- ctx.on('session/created', (session: Session) => {
- queue.push(frame({
- type: 'host/session-added',
- sessionId: session.id,
- // Derived at frame time like summarize(); a just-created session
- // has run no turn yet, so this is constantly true in practice.
- blank: sessionBlank(session),
- // Including cwd lets the client group the new session without refreshing the list.
- ...sessionListFields(session.header),
- }))
- }),
- ctx.on('session/disposed', (session: Session) => {
- queue.push(frame({ type: 'host/session-removed', sessionId: session.id }))
- }),
- ctx.on('agent/status', (agent: Agent, status: AgentStatus) => {
- queue.push(frame({ type: 'host/session-status', sessionId: agent.id, running: status === 'running' }))
- }),
- ctx.on('agent/error', (agent: Agent, _turn: number, _step: number, error: unknown) => {
- queue.push(frame({ type: 'host/agent-error', sessionId: agent.id, message: errorChain(error) }))
- }),
- ctx.on('domain/changed', (change) => {
- if (change.domain !== 'workspace') return
- if (change.table === '') {
- if (change.operation !== 'put') return
- const state = workspaceDomainState.parse(change.value)
- for (const workspaceId of state.workspaceIds) {
- if (committedWorkspaceIds.has(workspaceId)) continue
- const workspace = ctx.workspace.get(workspaceId)
- if (workspace === undefined) {
- throw new Error(`committed workspace registry references missing workspace "${workspaceId}"`)
- }
- committedWorkspaceIds.add(workspaceId)
- queue.push(frame({ type: 'host/workspace-changed', workspace: workspaceView(workspace) }))
- }
- if (state.archivedSessionIds.length !== archivedSessionIds.length
- || state.archivedSessionIds.some((id, index) => id !== archivedSessionIds[index])) {
- archivedSessionIds = state.archivedSessionIds
- queue.push(frame({
- type: 'host/archived-sessions-changed',
- archivedSessionIds: [...state.archivedSessionIds],
- }))
- }
- return
- }
- if (change.table !== 'workspaces') return
- if (change.operation === 'deleted') {
- if (!committedWorkspaceIds.delete(change.key)) return
- queue.push(frame({
- type: 'host/workspace-removed',
- workspaceId: change.key as WorkspaceId,
- }))
- return
- }
- if (!committedWorkspaceIds.has(change.key)) return
- // Existing-entity table writes are complete attach/touch commits.
- // A new entity's first put waits for the global registry write above.
- queue.push(frame({
- type: 'host/workspace-changed',
- workspace: changedWorkspaceView(change.key, change.value),
- }))
- }),
- ctx.on('commands/change', () => {
- queue.push(frame({ type: 'host/commands-changed' }))
- }),
- ctx.on('settings/document-updated', (ns) => {
- // The RAW-section event, not the resolved one: a field going from
- // inherited to overridden leaves the resolved value equal, and a
- // configuration client still has to re-read (its held revision is
- // stale, and the field's meaning changed).
- const name = String(ns)
- queue.push(frame({ type: 'host/settings-changed', ns: name }))
- // A provider's own settings carry its model catalog and endpoint,
- // so a change there invalidates the model list even when the route
- // set is untouched — `llm/adapters-updated` alone misses it.
- if (modelProviderNamespaces().has(name)) queue.push(frame({ type: 'host/models-changed' }))
- }),
- ctx.on('credentials/updated', (ref) => {
- queue.push(frame({ type: 'host/credentials-changed', ref: String(ref) }))
- }),
- ctx.on('llm/adapters-updated', () => {
- queue.push(frame({ type: 'host/models-changed' }))
- }),
- ]
- return queue.iterate(signal, () => { for (const dispose of disposers) dispose() })
- },
- },
- respond(message: ClientResponse): Promise<RpcReceipt> {
- // Route by the echoed rpcId (the wire correlation): approvals first,
- // then questions — the two registries share one id space of UUIDs.
- const approval = pendingApprovals.get(message.rpcId)
- if (approval !== undefined) {
- if (!message.result.ok) return Promise.resolve({ accepted: false, reason: 'bad-response' })
- const parsed = approvalResponsePayloadSchema.safeParse(message.result.value)
- // The payload's audit correlation must match the entry the rpcId routed
- // to — a mismatched answer is malformed, not merely late.
- if (!parsed.success || parsed.data.approvalId !== approval.approvalId || parsed.data.sessionId !== approval.sessionId) {
- return Promise.resolve({ accepted: false, reason: 'bad-response' })
- }
- approval.resolve(parsed.data.outcome)
- return Promise.resolve({ accepted: true })
- }
- const pending = pendingQuestions.get(message.rpcId)
- if (pending === undefined) return Promise.resolve({ accepted: false, reason: 'not-pending' })
- if (!message.result.ok) {
- if (message.result.error.code !== 'cancelled') {
- return Promise.resolve({ accepted: false, reason: 'bad-response' })
- }
- claimQuestion(pending, 'cancelled')
- pending.reject(new UserInteractionError(
- 'the user cancelled ask_user_question', 'ASK_CANCELLED'))
- return Promise.resolve({ accepted: true })
- }
- const parsed = questionResponsePayloadSchema.safeParse(message.result.value)
- if (!parsed.success) {
- return Promise.resolve({ accepted: false, reason: 'bad-response' })
- }
- const payload: QuestionResponsePayload = {
- sessionId: parsed.data.sessionId,
- answer: {
- answers: parsed.data.answer.answers.map(answer => ({
- id: answer.id,
- selected: answer.selected,
- ...(answer.custom === undefined ? {} : { custom: answer.custom }),
- })),
- },
- }
- if (!matchesQuestions(payload, pending)) {
- return Promise.resolve({ accepted: false, reason: 'bad-response' })
- }
- claimQuestion(pending, 'answered')
- pending.resolve(payload.answer)
- return Promise.resolve({ accepted: true })
- },
- }
- }
|