api-proxy.ts 165 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454345534563457345834593460346134623463346434653466346734683469347034713472347334743475347634773478347934803481348234833484348534863487348834893490349134923493349434953496349734983499350035013502350335043505350635073508350935103511351235133514351535163517351835193520352135223523352435253526352735283529353035313532353335343535353635373538353935403541354235433544354535463547354835493550355135523553355435553556355735583559356035613562356335643565356635673568356935703571357235733574357535763577357835793580358135823583358435853586358735883589359035913592359335943595359635973598359936003601360236033604360536063607360836093610361136123613361436153616361736183619362036213622362336243625362636273628362936303631363236333634363536363637363836393640364136423643364436453646364736483649365036513652365336543655365636573658365936603661366236633664366536663667366836693670367136723673367436753676367736783679368036813682368336843685368636873688368936903691369236933694369536963697369836993700370137023703370437053706370737083709371037113712371337143715371637173718371937203721372237233724372537263727372837293730373137323733373437353736373737383739374037413742374337443745374637473748374937503751375237533754375537563757375837593760376137623763376437653766376737683769377037713772377337743775377637773778377937803781378237833784378537863787378837893790379137923793379437953796379737983799380038013802380338043805380638073808380938103811381238133814381538163817381838193820382138223823382438253826382738283829383038313832383338343835383638373838383938403841
  1. /**
  2. * Host-side ApiProxy implementation. Signature discipline: unary takes the
  3. * narrow RpcRequest<P> and echoes request.rpcId on the RpcResponse<T>.
  4. */
  5. import { randomUUID } from 'node:crypto'
  6. import { mkdir, stat } from 'node:fs/promises'
  7. import { dirname } from 'node:path'
  8. import type { Context } from '@deepseek-ai/cordis'
  9. import { installModelSelection } from '@deepseek-ai/dsh-agent'
  10. import type { Agent, ModelSelection, ModelSelectionRef, AgentOptions, AgentStatus, PreStepDecision } from '@deepseek-ai/dsh-agent'
  11. import type {} from '@deepseek-ai/dsh-agent-presets/types'
  12. import { AttachmentError } from '@deepseek-ai/dsh-attachment'
  13. import type { ImageAttachmentRef } from '@deepseek-ai/dsh-attachment'
  14. import { contentHasImage, createUserMessage, freezeMessage, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  15. import { errorChain } from '@deepseek-ai/dsh-llm'
  16. import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
  17. import { isAppendSurfaceEvent, isJsonValue } from '@deepseek-ai/dsh-session'
  18. import type { JsonValue, Session, SessionEvent, SessionEventMap, SessionHeader, SessionId, UserMessage } from '@deepseek-ai/dsh-session'
  19. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  20. import {
  21. parseSessionReferenceText,
  22. type SessionReferenceInput,
  23. } from '@deepseek-ai/dsh-session-reference'
  24. import { SessionQueryError, type SessionSearchCursor } from '@deepseek-ai/dsh-session-query'
  25. import { SubagentError } from '@deepseek-ai/dsh-subagent'
  26. import type { SubagentListEntry as CatalogSubagentListEntry } from '@deepseek-ai/dsh-subagent'
  27. import { isUserInvocable } from '@deepseek-ai/dsh-skill'
  28. import type { Workspace, WorkspaceRecord } from '@deepseek-ai/dsh-workspace'
  29. import {
  30. workspaceDomainState, workspaceRecord, WorkspaceId as brandWorkspaceId,
  31. WorkspaceMoveInvalidError, WorkspaceOrderInvalidError, WorkspaceUnknownSessionError,
  32. } from '@deepseek-ai/dsh-workspace'
  33. // Type-only: brings the `ctx.tools` Context merge into this program (viewFor reads presenters).
  34. import {
  35. InvalidPresetIdError, PresetExistsError, PresetMountError,
  36. PresetNotWritableError, resolveSessionPreset, UnknownPresetError,
  37. } from '@deepseek-ai/dsh-agent-presets'
  38. import type { PresetBearingSession } from '@deepseek-ai/dsh-agent-presets'
  39. import type {} from '@deepseek-ai/dsh-tools'
  40. import type {
  41. ApiProxy, ConfigurableProviderView, CredentialView, GoalRef, HistoryEntry, HostFrame,
  42. ModelCatalogFailure, ModelProviderGroup,
  43. ModelReasoning, MuxFrame, PromptContentPart, QuestionResponsePayload, SessionListMetadata, SessionProjectionsBlock, SessionSearchItem,
  44. QueuedInboxItem, SessionSummary, SettingsNamespaceView, SubagentAddress, JobView, ToolEventView,
  45. WorkspaceId, WorkspaceView,
  46. } from './api/index.ts'
  47. import {
  48. DEFAULT_SESSION_LOG_COMPRESSION_LEVEL,
  49. flushLiveSessionLog,
  50. sessionLogExportDeps,
  51. sessionLogZipFilename,
  52. streamSessionLogZip,
  53. type SessionLogExportReady,
  54. type SessionLogCompressionLevel,
  55. } from './session-export.ts'
  56. import type { SessionRawArtifact } from '@deepseek-ai/dsh-session-persistence'
  57. import {
  58. SESSION_SEARCH_RESULT_LIMIT,
  59. SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS,
  60. truncateUnicodeCodePoints,
  61. } from './api/session-search.ts'
  62. // Type-only: resolves `ctx.get('sessionProjections')` to the projection registry.
  63. import type {} from '@deepseek-ai/dsh-session-projection'
  64. // Type-only: resolves `ctx.get('tasks')` to the background job registry.
  65. import type {} from '@deepseek-ai/dsh-jobs'
  66. import type { JobSnapshot } from '@deepseek-ai/dsh-jobs'
  67. // Type-only: resolves `ctx.get('sessionProjectionCache')` (the cold listing column).
  68. import type {} from '@deepseek-ai/dsh-session-projection-cache'
  69. // GoalError narrows domain rejections to their stable codes at the wire boundary.
  70. import { GoalError } from '@deepseek-ai/dsh-goal'
  71. import type { GoalRef as CoreGoalRef } from '@deepseek-ai/dsh-goal'
  72. // Type-only edges: resolve the command-change stream and `ctx.get('skills')`.
  73. import type {} from '@deepseek-ai/dsh-commands'
  74. // Type-only: the dynamic-package runner's forwarded-event declarations. Its
  75. // client-safe `./types` subpath deliberately, not the package root — the root
  76. // merges `ctx.dynamicCordisRunner`, and a dependency on that package would
  77. // rebuild the api-remotes cycle this direction exists to avoid.
  78. import type {} from '@deepseek-ai/dsh-cordis-host-runner/types'
  79. import type {} from '@deepseek-ai/dsh-skill'
  80. // The settings/credentials seams: brand guards run at this wire boundary; the
  81. // service reads stay optional (`ctx.get`) so a composition without either
  82. // provider still serves every other domain.
  83. import { SettingsConflictError, settingsNamespace } from '@deepseek-ai/dsh-settings'
  84. import type { SettingsDescriptor, SettingsNamespace, SettingsPathOp } from '@deepseek-ai/dsh-settings'
  85. import { credentialRef } from '@deepseek-ai/dsh-credentials'
  86. // Value edge: the rename impl narrows the title service's validation failure; the import also resolves `ctx.get('sessionTitle')`.
  87. import { SessionTitleInvalidError } from '@deepseek-ai/dsh-session-title'
  88. import type { CallId, MessageId } from '@deepseek-ai/dsh-llm/brand'
  89. import type { ScopeKey } from '@deepseek-ai/dsh-scope'
  90. import type { ApprovalOutcome, ApprovalRequestId } from '@deepseek-ai/dsh-user-approval'
  91. // Side-effect type import: resolves the `approval/request` waterfall and
  92. // `ctx.get('approval')` without a value dependency on the seam (optional composition).
  93. import type {} from '@deepseek-ai/dsh-user-approval'
  94. import { approvalResponsePayloadSchema } from './api/approvals.schema.ts'
  95. import { imageLimitsProjectionSchema, sessionListMetadataProjectionSchema } from './api/sessions.schema.ts'
  96. import { questionResponsePayloadSchema } from './api/questions.schema.ts'
  97. import type { ClientResponse, RpcError, RpcReceipt, RpcRequest, RpcResponse } from './api/rpc.ts'
  98. import { RpcId } from './api/rpc.ts'
  99. import type {
  100. AskUserQuestionAnswer, AskUserQuestionItem, AskUserQuestionRequest,
  101. } from '@deepseek-ai/dsh-user-questions'
  102. import { UserQuestionError } from '@deepseek-ai/dsh-user-questions'
  103. import { DirectoryPickerError } from '@deepseek-ai/dsh-host-directory-picker'
  104. import {
  105. ApiRemoteSessionNotFound as SessionNotFound,
  106. ApiRemoteSubagentSessionOwnership as SubagentSessionOwnership,
  107. API_REMOTE_FORWARDED_EVENTS,
  108. apiRemoteSubagentOwnershipError,
  109. createApiRemoteAgentResolver,
  110. hasApiRemoteSubagentOwner,
  111. inspectApiRemoteSession,
  112. } from '@deepseek-ai/dsh-api-remotes'
  113. import { canOpenNativePath, openNativePath, openNativeTextFile } from './native-path-opener.ts'
  114. /** Page size when history is called without maxMessages. */
  115. const DEFAULT_MAX_MESSAGES = 50
  116. /** Provider work budget: at most 100 calls and 2,000 inspected hits. */
  117. const SESSION_SEARCH_PROVIDER_CALL_LIMIT = 100
  118. /** Bound cold-log stat fan-out and settle each started batch before cancellation returns. */
  119. const COLD_SUMMARY_BATCH_SIZE = 16
  120. /** Default maximum artifact size eligible for one cold blankness read. */
  121. export const DEFAULT_COLD_BLANK_PROBE_MAX_BYTES = 1024
  122. /** Conversation message event types (the pagination counting unit). */
  123. const MESSAGE_TYPES = new Set(['user/message', 'assistant/message'])
  124. /** Decode the browser payload while rejecting non-canonical base64 forms. */
  125. function decodeBase64(data: string): Uint8Array {
  126. const decoded = Buffer.from(data, 'base64')
  127. if (data.length === 0 || decoded.toString('base64') !== data) {
  128. throw new AttachmentError('Image upload is not canonical base64.', 'INVALID_IMAGE_BASE64')
  129. }
  130. return new Uint8Array(decoded)
  131. }
  132. /** Validate one prompt as a batch before publishing any durable image object. */
  133. async function durablePromptContent(ctx: Context, content: readonly PromptContentPart[]): Promise<ContentBlock[]> {
  134. if (content.every(part => part.type === 'text')) {
  135. return content.map(part => ({ type: 'text', text: part.text }))
  136. }
  137. const limits = ctx.attachments.imageLimits
  138. if (content.filter(part => part.type === 'image').length > limits.maxImagesPerMessage) {
  139. throw new AttachmentError('Prompt exceeds the configured image-count limit.', 'TOO_MANY_IMAGES')
  140. }
  141. const prepared = content.map(part => part.type === 'text'
  142. ? part
  143. : { part, data: decodeBase64(part.data) })
  144. const images = prepared.filter((part): part is Extract<typeof part, { data: Uint8Array }> => 'data' in part)
  145. const totalBytes = images.reduce((sum, image) => sum + image.data.byteLength, 0)
  146. if (totalBytes > limits.maxMessageImageBytes) {
  147. throw new AttachmentError('Prompt exceeds the configured aggregate image-byte limit.', 'IMAGES_TOO_LARGE')
  148. }
  149. for (const image of images) {
  150. await ctx.attachments.validateImage({
  151. data: image.data,
  152. mediaType: image.part.mediaType,
  153. ...image.part.name === undefined ? {} : { name: image.part.name },
  154. })
  155. }
  156. const blocks: ContentBlock[] = []
  157. for (const item of prepared) {
  158. if (!('data' in item)) {
  159. blocks.push({ type: 'text', text: item.text })
  160. continue
  161. }
  162. const attachment = await ctx.attachments.saveImage({
  163. data: item.data,
  164. mediaType: item.part.mediaType,
  165. ...item.part.name === undefined ? {} : { name: item.part.name },
  166. })
  167. blocks.push({ type: 'image', attachment })
  168. }
  169. return blocks
  170. }
  171. /** Remove canonical session mentions from text blocks and retain their structured identities. */
  172. function parseReferencedContent(content: readonly PromptContentPart[]): {
  173. content: PromptContentPart[]
  174. references: SessionReferenceInput[]
  175. } {
  176. const references: SessionReferenceInput[] = []
  177. const normalized = content.map((part): PromptContentPart => {
  178. if (part.type !== 'text') return part
  179. const parsed = parseSessionReferenceText(part.text)
  180. references.push(...parsed.references)
  181. return { type: 'text', text: parsed.text }
  182. })
  183. return { content: normalized, references }
  184. }
  185. /** Search durable content for an image reference, including nested tool results. */
  186. function imageBlockIn(content: unknown, match: (ref: ImageAttachmentRef) => boolean): ImageAttachmentRef | undefined {
  187. if (!Array.isArray(content)) return undefined
  188. for (const value of content) {
  189. if (typeof value !== 'object' || value === null || Array.isArray(value)) continue
  190. const block = value as { type?: unknown; attachment?: unknown; content?: unknown }
  191. if (block.type === 'image' && typeof block.attachment === 'object' && block.attachment !== null) {
  192. const ref = block.attachment as ImageAttachmentRef
  193. if (match(ref)) return ref
  194. }
  195. if (block.type === 'tool-result') {
  196. const nested = imageBlockIn(block.content, match)
  197. if (nested !== undefined) return nested
  198. }
  199. }
  200. return undefined
  201. }
  202. /** Search every durable event carrier that can own model-visible content. */
  203. function imageInEvent(event: SessionEvent, match: (ref: ImageAttachmentRef) => boolean): ImageAttachmentRef | undefined {
  204. const data = event.data as {
  205. content?: unknown
  206. message?: { content?: unknown }
  207. inserted?: Array<{ content?: unknown }>
  208. chunk?: { type?: unknown; block?: unknown }
  209. }
  210. const direct = imageBlockIn(data.content, match)
  211. if (direct !== undefined) return direct
  212. if (data.message !== undefined) {
  213. const wrapped = imageBlockIn(data.message.content, match)
  214. if (wrapped !== undefined) return wrapped
  215. }
  216. if (data.inserted !== undefined) {
  217. for (const message of data.inserted) {
  218. const inserted = imageBlockIn(message.content, match)
  219. if (inserted !== undefined) return inserted
  220. }
  221. }
  222. if (event.type === 'assistant/chunk' && data.chunk?.type === 'block-end') {
  223. return imageBlockIn([data.chunk.block], match)
  224. }
  225. return undefined
  226. }
  227. /** True when the current model-visible surface contains an image. */
  228. function messagesHaveImage(messages: readonly { content: readonly ContentBlock[] }[]): boolean {
  229. return messages.some(message => contentHasImage(message.content))
  230. }
  231. /** Resolve the first reference matching one opaque id. */
  232. function referencedImage(events: readonly SessionEvent[], attachmentId: string): ImageAttachmentRef | undefined {
  233. for (const event of events) {
  234. const found = imageInEvent(event, ref => String(ref.attachmentId) === attachmentId)
  235. if (found !== undefined) return found
  236. }
  237. return undefined
  238. }
  239. /** Strict browser-zone profile: UTC or an IANA Area/Location-style identifier. */
  240. const IANA_TIME_ZONE = /^[A-Za-z][A-Za-z0-9_+.-]*(?:\/[A-Za-z0-9_+.-]+)+$/
  241. /** Validate and canonicalize one browser-supplied IANA zone at the wire boundary. */
  242. function canonicalClientTimeZone(value: string): string | undefined {
  243. if (value.length === 0 || value.trim() !== value
  244. || (value !== 'UTC' && !IANA_TIME_ZONE.test(value))) return undefined
  245. try {
  246. const canonical = new Intl.DateTimeFormat('en-US', { timeZone: value })
  247. .resolvedOptions().timeZone
  248. /* v8 ignore next -- Intl returns UTC or a canonical IANA Area/Location for accepted input. */
  249. if (canonical !== 'UTC' && !IANA_TIME_ZONE.test(canonical)) return undefined
  250. return canonical
  251. } catch {
  252. // Intl rejects unsupported zone names; the RPC maps that parser rejection below.
  253. return undefined
  254. }
  255. }
  256. /** Read live abort state across awaits without treating it as synchronously immutable. */
  257. function isAborted(signal: AbortSignal): boolean {
  258. return signal.aborted
  259. }
  260. /**
  261. * Message-boundary pagination: count maxMessages append-origin messages
  262. * backwards from the window tail. Replacement copies never entered the
  263. * conversation a reader sees — they restate a shadowed range for the model
  264. * alone — so they consume no quota; the page stays one contiguous raw range,
  265. * which keeps a compaction's log-only `compaction/summary` record on the same page as its
  266. * replacement. The cut is the starting seq of the oldest message group (chunks
  267. * group via sourceEventSeqs — never cut mid-message). The tail page naturally
  268. * includes the in-progress partial.
  269. */
  270. function paginate(
  271. events: readonly SessionEvent[],
  272. beforeSeq: number | undefined,
  273. maxMessages: number,
  274. ): { events: SessionEvent[]; hasMore: boolean } {
  275. const window = beforeSeq === undefined ? [...events] : events.filter(event => event.seq < beforeSeq)
  276. let count = 0
  277. let cut = 0
  278. for (let i = window.length - 1; i >= 0; i--) {
  279. const event = window[i] as SessionEvent
  280. if (!MESSAGE_TYPES.has(event.type) || !isAppendSurfaceEvent(event)) continue
  281. count++
  282. const sources = (event as { sourceEventSeqs?: number[] }).sourceEventSeqs
  283. let groupStart = event.seq
  284. if (sources !== undefined) {
  285. for (const source of sources) {
  286. if (source < groupStart) groupStart = source
  287. }
  288. }
  289. if (count >= maxMessages) {
  290. cut = groupStart
  291. break
  292. }
  293. }
  294. const page = window.filter(event => event.seq >= cut)
  295. return { events: page, hasMore: cut > 0 }
  296. }
  297. /** Wrap an ok result echoing the request's rpcId. */
  298. function ok<T>(request: RpcRequest<unknown>, value: T): RpcResponse<T> {
  299. return { rpcId: request.rpcId, result: { ok: true, value } }
  300. }
  301. /**
  302. * Build the provider/model catalog over every registered route. Shared by the
  303. * session-scoped `session.models` and host-scoped `llm.models`. Catalog
  304. * membership stays advisory: an unlisted session selection remains valid for
  305. * provider dispatch, but is not injected back into the selector after its
  306. * owning catalog stops advertising it. Per-provider failures ride `failures`
  307. * without failing the sound groups; groups that advertise nothing are dropped.
  308. */
  309. async function buildModelCatalog(ctx: Context): Promise<{
  310. groups: ModelProviderGroup[]
  311. failures: ModelCatalogFailure[]
  312. }> {
  313. const catalog = await Promise.all(ctx.llm.listProviders().map(async (provider) => {
  314. try {
  315. const models = await ctx.llm.listModels(provider.id)
  316. const entries = await Promise.all(models.map(async (model) => {
  317. const resolved = await ctx.llm.resolveModelInfo(provider.id, model.id)
  318. const reasoning: ModelReasoning | undefined = resolved.reasoning === undefined
  319. ? undefined
  320. : {
  321. efforts: resolved.reasoning.efforts.map(effort => ({
  322. id: effort.id,
  323. name: effort.name,
  324. ...effort.description === undefined
  325. ? {}
  326. : { description: effort.description },
  327. })),
  328. ...resolved.reasoning.defaultEffort === undefined
  329. ? {}
  330. : { defaultEffort: resolved.reasoning.defaultEffort },
  331. }
  332. return {
  333. id: model.id,
  334. name: model.name,
  335. ...model.description === undefined ? {} : { description: model.description },
  336. ...reasoning === undefined ? {} : { reasoning },
  337. }
  338. }))
  339. const group: ModelProviderGroup = {
  340. id: provider.id,
  341. name: provider.name,
  342. models: entries,
  343. }
  344. return { kind: 'group' as const, group }
  345. } catch (error: unknown) {
  346. const failure: ModelCatalogFailure = {
  347. id: provider.id,
  348. name: provider.name,
  349. message: error instanceof Error ? error.message : String(error),
  350. }
  351. return { kind: 'failure' as const, failure }
  352. }
  353. }))
  354. return {
  355. groups: catalog.flatMap(item => item.kind === 'group' ? [item.group] : []).filter(group => group.models.length > 0),
  356. failures: catalog.flatMap(item => item.kind === 'failure' ? [item.failure] : []),
  357. }
  358. }
  359. /** Wrap an error result echoing the request's rpcId. */
  360. function err<T>(request: RpcRequest<unknown>, error: RpcError): RpcResponse<T> {
  361. return { rpcId: request.rpcId, result: { ok: false, error } }
  362. }
  363. /**
  364. * The RPC refusal a preset failure becomes, or undefined when the failure is
  365. * about something else.
  366. *
  367. * Both the session-create path and the switch path can be handed the same two
  368. * failures, and a client that has to branch on the code needs them worded the
  369. * same from either.
  370. * @param request - the request being answered.
  371. * @param error - the thrown value.
  372. * @returns the refusal, or undefined when the caller should keep handling.
  373. */
  374. function presetFailure(request: RpcRequest<unknown>, error: unknown): RpcResponse<never> | undefined {
  375. if (error instanceof UnknownPresetError) {
  376. return err(request, {
  377. code: 'agent-preset-not-found',
  378. message: error.message,
  379. details: { agentPreset: error.presetId, available: [...error.available] },
  380. })
  381. }
  382. if (error instanceof PresetMountError) {
  383. return err(request, {
  384. code: 'agent-preset-invalid',
  385. message: error.message,
  386. details: { agentPreset: error.presetId, reason: error.reason },
  387. })
  388. }
  389. return undefined
  390. }
  391. /** Simple async queue: core callbacks push, the AsyncIterable pulls; abort/return cleans up. */
  392. class FrameQueue<F> {
  393. private buffer: F[] = []
  394. private waiter: (() => void) | undefined
  395. private done = false
  396. push(item: F): void {
  397. if (this.done) return
  398. this.buffer.push(item)
  399. this.waiter?.()
  400. }
  401. end(): void {
  402. this.done = true
  403. this.waiter?.()
  404. }
  405. async *iterate(signal: AbortSignal, cleanup: () => void): AsyncGenerator<F> {
  406. const onAbort = (): void => { this.end() }
  407. signal.addEventListener('abort', onAbort, { once: true })
  408. try {
  409. while (true) {
  410. while (this.buffer.length > 0) yield this.buffer.shift() as F
  411. if (this.done || signal.aborted) return
  412. await new Promise<void>((resolve) => { this.waiter = resolve })
  413. this.waiter = undefined
  414. }
  415. } finally {
  416. signal.removeEventListener('abort', onAbort)
  417. cleanup()
  418. }
  419. }
  420. }
  421. /**
  422. * Server-side frame mint: pure pushes get a fresh rpcId per frame (answerable
  423. * frames — approval/question requested — mint their stable id in their
  424. * pending registries instead).
  425. */
  426. function frame<F>(payload: F): RpcRequest<F> {
  427. return { rpcId: RpcId(randomUUID()), payload }
  428. }
  429. /**
  430. * Narrow one allowlisted host event's argument list to the JSON values the
  431. * wrapper frame carries. A rejected argument is an allowlist mistake (the
  432. * forwarded path applies no projection), not hostile input, so it throws rather
  433. * than degrading to a lossy frame. The throw surfaces where the forwarding
  434. * listener runs, so the emitter's own listener containment logs it and drops
  435. * that frame — loud in the Host log, not at load or at the emit. Exported for
  436. * the test that owns this decision: every currently allowlisted event has a
  437. * statically JSON-safe payload, so a type-legal `ctx.emit` cannot reach the
  438. * rejection branch.
  439. * @param event - forwarded host event name, named in the failure.
  440. * @param args - the emitter's argument list.
  441. * @returns the same arguments typed as JSON values.
  442. */
  443. export function assertJsonArgs(event: string, args: readonly unknown[]): JsonValue[] {
  444. for (const [index, arg] of args.entries()) {
  445. if (!isJsonValue(arg)) {
  446. throw new Error(`forwarded host event "${event}" argument ${index} is not lossless JSON data`)
  447. }
  448. }
  449. return args as JsonValue[]
  450. }
  451. /** Queue the subscription baseline frame. */
  452. function subscribeSession(queue: FrameQueue<RpcRequest<MuxFrame>>, session: Session): void {
  453. queue.push(frame({ type: 'session/subscribed', sessionId: session.id, lastSeq: session.seq - 1 }))
  454. }
  455. /**
  456. * Project registry snapshots onto the wire view, dropping the three internal
  457. * fields {@link JobView} documents as absent.
  458. */
  459. function jobViews(snapshots: readonly JobSnapshot[]): JobView[] {
  460. return snapshots.map(job => ({
  461. id: job.id,
  462. kind: job.kind,
  463. label: job.label,
  464. status: job.status,
  465. ...job.detail === undefined ? {} : { detail: job.detail },
  466. startedAt: job.startedAt,
  467. ...job.finishedAt === undefined ? {} : { finishedAt: job.finishedAt },
  468. }))
  469. }
  470. /**
  471. * Whether the session's conversation has started: no turn has run yet (a
  472. * turn is one model-loop execution). Standalone plugin events — command
  473. * lifecycle records, plan/mode, titles, goals — never open a turn, so
  474. * running `/plan` or `/goal` on a fresh session keeps it blank
  475. * (list-hidden, reusable).
  476. */
  477. function sessionBlank(session: Session): boolean {
  478. return !session.events.some(event => event.type === 'turn/start')
  479. }
  480. /** Advance the Session-list hint projection by one committed event. */
  481. function applySessionListMetadata(state: SessionListMetadata, event: SessionEvent): SessionListMetadata {
  482. const blank = state.blank && event.type !== 'turn/start'
  483. const lastPromptAt = event.type === 'user/message' && event.data.source.kind === 'user'
  484. ? event.time
  485. : state.lastPromptAt
  486. return blank === state.blank && lastPromptAt === state.lastPromptAt
  487. ? state
  488. : { blank, lastPromptAt }
  489. }
  490. /** Fold exact list metadata for an attached Session. */
  491. function sessionListMetadata(events: readonly SessionEvent[]): SessionListMetadata {
  492. let state: SessionListMetadata = { blank: true, lastPromptAt: null }
  493. for (const event of events) state = applySessionListMetadata(state, event)
  494. return state
  495. }
  496. /** Sort by creation or latest human prompt, whichever is newer. */
  497. function sessionListUpdatedAt(header: SessionHeader, metadata: SessionListMetadata | undefined): number {
  498. return Math.max(header.createdAt, metadata?.lastPromptAt ?? 0)
  499. }
  500. /** Shared Session-header projection for list baselines and creation frames. */
  501. function sessionListFields(header: SessionHeader, events: readonly SessionEvent[] = []): {
  502. parentSessionId?: SessionId
  503. origin?: 'subagent'
  504. cwd?: string
  505. agentPreset?: string
  506. } {
  507. // The preset comes from the log, not the header: a session that switched
  508. // while blank ran its turns under the newer composition, and a picker
  509. // showing the creation-time value would contradict what the model saw.
  510. const agentPreset = resolveSessionPreset({ header, events })
  511. return {
  512. ...header.parentSession === undefined ? {} : { parentSessionId: header.parentSession },
  513. ...header.origin === undefined ? {} : { origin: header.origin },
  514. ...header.cwd === undefined ? {} : { cwd: header.cwd },
  515. ...agentPreset === undefined ? {} : { agentPreset },
  516. }
  517. }
  518. /** SessionSummary projection for attached (in-memory) sessions. */
  519. function summarize(session: Session, running: boolean): SessionSummary {
  520. const metadata = sessionListMetadata(session.events)
  521. return {
  522. sessionId: session.id,
  523. updatedAt: sessionListUpdatedAt(session.header, metadata),
  524. running,
  525. blank: metadata.blank,
  526. ...sessionListFields(session.header, session.events),
  527. }
  528. }
  529. /**
  530. * Verify a possibly blank cold Session only when its physical artifact passes
  531. * the configured per-Session size check. A stale `blank: true`, an
  532. * absent cache row, a large or location-less artifact, and read failures all
  533. * resolve to visible (`false`); listing must never hide a conversation on a
  534. * cache hint or an unavailable optimization.
  535. */
  536. async function probeColdSessionMetadata(
  537. ctx: Context,
  538. persistence: SessionPersistence,
  539. meta: SessionHeader,
  540. maxBytes: number,
  541. signal?: AbortSignal,
  542. ): Promise<SessionListMetadata | undefined> {
  543. if (maxBytes === 0) return undefined
  544. signal?.throwIfAborted()
  545. const location = persistence.locate(meta)
  546. if (location === undefined) return undefined
  547. signal?.throwIfAborted()
  548. let size: number
  549. try {
  550. size = (await stat(location.path)).size
  551. } catch {
  552. signal?.throwIfAborted()
  553. return undefined
  554. }
  555. if (size > maxBytes) return undefined
  556. try {
  557. const { events } = await persistence.readFrom(meta.id, 0, signal)
  558. signal?.throwIfAborted()
  559. return sessionListMetadata(events)
  560. } catch (error) {
  561. signal?.throwIfAborted()
  562. ctx.logger.warn(`session.list: blank probe for "${meta.id}" failed (serving it as visible): ${String(error)}`)
  563. return undefined
  564. }
  565. }
  566. /** SessionSummary projection for a cold persisted Session. */
  567. async function summarizeCold(
  568. ctx: Context,
  569. persistence: SessionPersistence,
  570. meta: SessionHeader,
  571. metadata: SessionListMetadata | undefined,
  572. blankProbeMaxBytes: number,
  573. signal?: AbortSignal,
  574. ): Promise<SessionSummary> {
  575. const probed = metadata?.blank === false
  576. ? undefined
  577. : await probeColdSessionMetadata(ctx, persistence, meta, blankProbeMaxBytes, signal)
  578. return {
  579. sessionId: meta.id,
  580. updatedAt: sessionListUpdatedAt(meta, probed ?? metadata),
  581. running: false,
  582. blank: metadata?.blank === false ? false : probed?.blank ?? false,
  583. // Header-only: reading the log for a blank-window preset switch would
  584. // defeat the same index read, and attaching the session replaces this row
  585. // with `summarize()`, which resolves the switch from the events.
  586. ...sessionListFields(meta),
  587. }
  588. }
  589. /** Map a browse-primitive failure onto the wire error vocabulary (unknown throws stay internal). */
  590. function directoryError(error: unknown): RpcError {
  591. if (error instanceof DirectoryPickerError) {
  592. return { code: error.code, message: error.message, details: { path: error.path } }
  593. }
  594. return { code: 'internal', message: error instanceof Error ? error.message : String(error), details: {} }
  595. }
  596. /** Resolved Agent model and project-directory defaults consumed by the API implementation. */
  597. export interface ApiProxyDefaults {
  598. /**
  599. * The model selection a session starts from when its own log names none. Read on
  600. * every access rather than captured, so a default saved during this process
  601. * reaches the sessions that have not run a turn yet.
  602. */
  603. defaultModelSelection: () => ModelSelection
  604. /**
  605. * Record a selection as the new default. Either absent, or a closure that
  606. * may itself decline — the gateway plugin always passes one, and it no-ops
  607. * when the deployment mounts no settings provider or when the write races
  608. * service teardown. A switch then stays process-local. A rejection is
  609. * reported and swallowed: the switch already applies to its own session,
  610. * and undoing it because storage failed would be the worse outcome.
  611. */
  612. saveDefaultModelSelection?: (selection: ModelSelection) => Promise<void>
  613. /** Default project directory for new sessions whose create request carries no cwd. */
  614. cwd: string
  615. /** Native open-with-default-application; injectable for carrier tests. */
  616. openPath?: (path: string, signal: AbortSignal) => Promise<void>
  617. /** Native text-editor handoff; injectable for settings-document tests. */
  618. openTextFile?: (path: string, signal: AbortSignal) => Promise<void>
  619. /** Validated DEFLATE level for session-log ZIP entries; defaults to 6. */
  620. sessionExportCompressionLevel?: SessionLogCompressionLevel
  621. /** Maximum artifact size eligible for one cold blankness read. */
  622. coldBlankProbeMaxBytes?: number
  623. /**
  624. * Whether handing a path to the native opener can work at all — the
  625. * `hasDocument` capability the preset roster reports, and the switch
  626. * between opening a preset directory and answering its path as text.
  627. * Absent, an injected `openPath` counts as openable and everything else
  628. * falls back to platform detection ({@link canOpenNativePath}).
  629. */
  630. canOpenPath?: () => boolean
  631. }
  632. /** The tool/call payload fields the presenter path reads. */
  633. interface ToolCallData { callId: string; name: string; arguments: string }
  634. /**
  635. * One outstanding approval question: the stable server-request id, the frame
  636. * material replayed to late mux subscribers, and the resolver that settles the
  637. * answerer's promise back into `ctx.approval`.
  638. */
  639. interface PendingApproval {
  640. rpcId: RpcId
  641. sessionId: SessionId
  642. approvalId: ApprovalRequestId
  643. toolName: string
  644. callId?: CallId
  645. reason?: string
  646. resolve(outcome: ApprovalOutcome): void
  647. }
  648. /** Project a pending entry into its answerable mux frame (initial push and mux-open replay share it). */
  649. function requestedFrame(pending: PendingApproval): RpcRequest<MuxFrame> {
  650. return {
  651. rpcId: pending.rpcId,
  652. payload: {
  653. type: 'approval/requested',
  654. sessionId: pending.sessionId,
  655. approvalId: pending.approvalId,
  656. toolName: pending.toolName,
  657. ...pending.callId === undefined ? {} : { callId: pending.callId },
  658. ...pending.reason === undefined ? {} : { reason: pending.reason },
  659. },
  660. }
  661. }
  662. /** One host-owned question wait, addressed by the stable server-request id. */
  663. interface PendingQuestion {
  664. rpcId: RpcId
  665. sessionId: SessionId
  666. questions: AskUserQuestionItem[]
  667. resolve: (answer: AskUserQuestionAnswer) => void
  668. reject: (error: UserQuestionError) => void
  669. signal?: AbortSignal
  670. onAbort?: () => void
  671. }
  672. /** Validate one answer batch against the exact question request it resolves. */
  673. function matchesQuestions(payload: QuestionResponsePayload, pending: PendingQuestion): boolean {
  674. if (payload.sessionId !== pending.sessionId) return false
  675. const answers = payload.answer.answers
  676. if (answers.length !== pending.questions.length) return false
  677. return answers.every((answer, index) => {
  678. const question = pending.questions[index] as AskUserQuestionItem
  679. if (answer.id !== question.id) return false
  680. if (new Set(answer.selected).size !== answer.selected.length) return false
  681. const custom = answer.custom?.trim()
  682. if (custom !== undefined && custom === '') return false
  683. if (question.multiSelect !== true) {
  684. if (custom !== undefined && answer.selected.length > 0) return false
  685. if (answer.selected.length > 1) return false
  686. }
  687. const labels = new Set(question.options?.map(option => option.label) ?? [])
  688. return answer.selected.every(label => labels.has(label))
  689. })
  690. }
  691. /**
  692. * Compute the render intent for a tool/call or tool/result event through the
  693. * presenters registered at this moment; every other event type gets none. A
  694. * result's presenter needs its call's parsed args — `argsFor` supplies them
  695. * (live: the per-session call table; history: an in-page backscan), returning
  696. * undefined when the pairing is unavailable (e.g. the call fell off the page),
  697. * which soft-falls to no view. Presenter or JSON.parse throws also soft-fall:
  698. * the client's documented default (generic JSON card) covers every miss.
  699. */
  700. function viewFor(
  701. ctx: Context,
  702. event: SessionEvent,
  703. argsFor: (callId: string) => unknown,
  704. // Presenters live with the definitions, and definitions live in the scope
  705. // chain: a preset registers its tools into its standing layer. A live agent
  706. // is a scope whose chain passes through its preset; a cold read passes the
  707. // preset's standing key directly — no agent, no resume. An undefined scope
  708. // sees only the global layer, which is the pre-preset deployment shape.
  709. scope?: ScopeKey,
  710. ): ToolEventView | undefined {
  711. try {
  712. if (event.type === 'tool/call') {
  713. const { name, arguments: raw } = event.data as ToolCallData
  714. const view = ctx.tools.get(name, scope)?.presentCall?.(JSON.parse(raw))
  715. return view === undefined ? undefined : { for: 'call', view }
  716. }
  717. if (event.type === 'tool/result') {
  718. const { message, meta } = event.data
  719. const [result] = message.content
  720. const callId = message.source.callId
  721. const call = argsFor(callId) as { name: string; args: unknown } | undefined
  722. if (call === undefined) return undefined
  723. const view = ctx.tools.get(call.name, scope)?.presentResult?.(call.args, {
  724. content: result.content,
  725. isError: result.isError === true,
  726. ...meta === undefined ? {} : { meta },
  727. })
  728. return view === undefined ? undefined : { for: 'result', view }
  729. }
  730. } catch (error: unknown) {
  731. // A throwing presenter (or unparseable arguments) must not break delivery;
  732. // the event still ships, just without a view.
  733. console.error(`api-proxy: presenter failed for ${event.type}, falling back to generic: ${String(error)}`)
  734. }
  735. return undefined
  736. }
  737. /**
  738. * Resolve a tool/result's call pairing by scanning a window of events backwards
  739. * for the matching tool/call. Used by the history path (the page is the
  740. * window — a cross-page pairing soft-falls to no view) and by live-path table
  741. * misses after a reconnect-eviction.
  742. */
  743. function backscanArgs(events: readonly SessionEvent[], callId: string): { name: string; args: unknown } | undefined {
  744. for (let i = events.length - 1; i >= 0; i--) {
  745. const event = events[i] as SessionEvent
  746. if (event.type !== 'tool/call') continue
  747. const data = event.data as ToolCallData
  748. if (data.callId !== callId) continue
  749. try {
  750. return { name: data.name, args: JSON.parse(data.arguments) }
  751. } catch {
  752. // Unparseable stored arguments: same soft-fall as a live parse failure.
  753. return undefined
  754. }
  755. }
  756. return undefined
  757. }
  758. /** Render one detached history page through the same presenter path as ordinary history. */
  759. function historyPage(
  760. ctx: Context,
  761. events: readonly SessionEvent[],
  762. beforeSeq: number | undefined,
  763. maxMessages: number | undefined,
  764. scope?: ScopeKey,
  765. ): { events: HistoryEntry[]; hasMore: boolean } {
  766. const page = paginate(events, beforeSeq, maxMessages ?? DEFAULT_MAX_MESSAGES)
  767. return {
  768. events: page.events.map((event) => {
  769. const view = viewFor(ctx, event, callId => backscanArgs(page.events, callId), scope)
  770. return { event, ...view === undefined ? {} : { view } }
  771. }),
  772. hasMore: page.hasMore,
  773. }
  774. }
  775. /**
  776. * The projection baseline for one history tail page: the registry's
  777. * watermark-cache snapshot — one fully synchronous read (no await between the
  778. * page slice and this), so all values and `asOfSeq` form a single consistent
  779. * cut and `asOfSeq` equals the window tail event seq. The carrier holds zero
  780. * domain knowledge (each value passed its unit's own schema inside the
  781. * registry). An absent registry means the deployment has no projection seam:
  782. * the whole block is absent and clients treat every key as capability-absent.
  783. */
  784. /**
  785. * Which session a transcript read is served from. An attached session is the
  786. * live object and keeps appending, so its events and projection baseline are
  787. * read together in one synchronous step; a detached one is already a frozen
  788. * inspection.
  789. */
  790. type HistorySource =
  791. | { readonly kind: 'attached'; readonly session: Session }
  792. | { readonly kind: 'detached'; readonly header: SessionHeader; readonly events: SessionEvent[] }
  793. function projectionsFor(ctx: Context, session: Session): SessionProjectionsBlock | undefined {
  794. const registry = ctx.get('sessionProjections')
  795. if (registry === undefined) return undefined
  796. return registry.snapshot(session)
  797. }
  798. /**
  799. * The projection baseline of one session.list row, fail-soft: attached
  800. * sessions cut the registry's live watermark cache; cold sessions view the
  801. * persisted projection cache's identity-checked stored rows (zero log loads
  802. * either way — the listing use case the cache exists for). The block shape
  803. * (values + asOfSeq) matches the history tail's, so a client seeds its
  804. * value store under the same higher-seq-wins rule. Any failure — and an
  805. * empty value set — yields an absent block: a listing without projections
  806. * is degraded, never broken.
  807. */
  808. function listProjectionsFor(ctx: Context, meta: SessionHeader, session: Session | undefined): SessionProjectionsBlock | undefined {
  809. try {
  810. const block = session !== undefined
  811. ? ctx.get('sessionProjections')?.snapshot(session)
  812. : ctx.get('sessionProjectionCache')?.cachedSnapshot(meta)
  813. return block !== undefined && Object.keys(block.values).length > 0 ? block : undefined
  814. } catch (error) {
  815. ctx.logger.warn(`session.list: projection column for "${meta.id}" failed (serving the row without it): ${String(error)}`)
  816. return undefined
  817. }
  818. }
  819. /** Projection baseline for a detached history tail without Agent activation. */
  820. function detachedProjectionsFor(
  821. ctx: Context,
  822. events: readonly SessionEvent[],
  823. ): SessionProjectionsBlock | undefined {
  824. const registry = ctx.get('sessionProjections')
  825. if (registry === undefined) return undefined
  826. return registry.restore({}, events, 0).snapshot
  827. }
  828. /**
  829. * Best-effort projections for one subagent history page, fail-soft like
  830. * {@link listProjectionsFor}: a registered unit throwing on a corrupt payload
  831. * never blocks transcript reading — the page is served without the block.
  832. * @param ctx - context carrying the logger for the degradation warning.
  833. * @param childSessionId - the child whose page is being decorated.
  834. * @param compute - the arm-specific fold (live watermark or detached restore).
  835. * @returns the projections block, or undefined when the fold failed.
  836. */
  837. function subagentHistoryProjections(
  838. ctx: Context,
  839. childSessionId: SessionId,
  840. compute: () => SessionProjectionsBlock | undefined,
  841. ): SessionProjectionsBlock | undefined {
  842. try {
  843. return compute()
  844. } catch (error) {
  845. ctx.logger.warn(`subagent.history: projections for "${childSessionId}" failed (serving the page without them): ${String(error)}`)
  846. return undefined
  847. }
  848. }
  849. /** Map continuation admission failures without exposing provider details. */
  850. function subagentPromptError(
  851. request: RpcRequest<{ childSessionId: SessionId }>,
  852. error: unknown,
  853. signal: AbortSignal,
  854. ): RpcResponse<never> {
  855. const childSessionId = request.payload.childSessionId
  856. if (signal.aborted) {
  857. return err(request, { code: 'cancelled', message: 'subagent prompt was cancelled', details: {} })
  858. }
  859. if (error instanceof SubagentError) {
  860. switch (error.code) {
  861. case 'NOT_RESUMABLE':
  862. return err(request, {
  863. code: 'subagent-not-resumable',
  864. message: 'subagent cannot be resumed',
  865. details: { childSessionId },
  866. })
  867. case 'UNAUTHORIZED':
  868. return err(request, {
  869. code: 'subagent-unauthorized',
  870. message: 'subagent does not belong to this parent',
  871. details: { childSessionId },
  872. })
  873. case 'DRAINING':
  874. case 'ACTIVATION_CLOSING':
  875. case 'CONTINUATION_UNAVAILABLE':
  876. case 'PERSISTENCE_UNAVAILABLE':
  877. return err(request, {
  878. code: 'subagent-delivery-unavailable',
  879. message: 'subagent follow-up is temporarily unavailable',
  880. details: { childSessionId },
  881. })
  882. default:
  883. break
  884. }
  885. }
  886. return err(request, { code: 'internal', message: 'subagent prompt failed', details: {} })
  887. }
  888. /** Stable RPC face of the missing projections capability, shared by every catalog read path. */
  889. function projectionsUnavailableError(): RpcError {
  890. return {
  891. code: 'internal',
  892. message: 'subagent catalog is unavailable: this deployment does not mount the sessionProjections registry (load @deepseek-ai/dsh-session-projection)',
  893. details: {},
  894. }
  895. }
  896. /** Verify one address and mode against the complete direct-child catalog. */
  897. async function catalogChild(
  898. ctx: Context,
  899. address: SubagentAddress,
  900. signal?: AbortSignal,
  901. ): Promise<{
  902. entry?: Extract<CatalogSubagentListEntry, { kind: 'child' }>
  903. error?: RpcError
  904. }> {
  905. const { parentSessionId, childSessionId, mode } = address
  906. try {
  907. const entries = await ctx.subagents.listChildren(parentSessionId, signal)
  908. const entry = entries.find(candidate => candidate.id === childSessionId)
  909. if (entry === undefined || (entry.kind === 'child' && entry.mode !== mode)) {
  910. return {
  911. error: {
  912. code: 'subagent-not-found',
  913. message: `session "${childSessionId}" is not a ${mode} direct child of "${parentSessionId}"`,
  914. details: { parentSessionId, childSessionId },
  915. },
  916. }
  917. }
  918. if (entry.kind === 'diagnostic') {
  919. return {
  920. error: {
  921. code: 'subagent-catalog-diagnostic',
  922. message: `subagent "${childSessionId}" is ${entry.reason}`,
  923. details: { parentSessionId, childSessionId, reason: entry.reason },
  924. },
  925. }
  926. }
  927. return { entry }
  928. } catch (error: unknown) {
  929. if (signal?.aborted || (error instanceof SubagentError && error.code === 'CANCELLED')) {
  930. return { error: { code: 'cancelled', message: 'subagent catalog read was cancelled', details: {} } }
  931. }
  932. if (error instanceof SubagentError && error.code === 'SUBAGENT_CONTROL_PROJECTIONS_UNAVAILABLE') {
  933. return { error: projectionsUnavailableError() }
  934. }
  935. return { error: { code: 'internal', message: 'subagent catalog read failed', details: {} } }
  936. }
  937. }
  938. /**
  939. * The requested preset differs from the one this session already runs.
  940. *
  941. * A session's composition is fixed at creation: its history was produced under
  942. * that preset's tools, so adopting the identity under a different one would
  943. * replay tool calls the rebuilt agent cannot make. Naming a different preset
  944. * is therefore a caller error rather than a switch.
  945. */
  946. /** The roster is absent: this deployment composes no agent presets at all. */
  947. function noRoster(agentPreset: string): RpcError {
  948. return {
  949. code: 'agent-preset-not-found',
  950. message: 'this deployment composes no agent presets',
  951. details: { agentPreset, available: [] },
  952. }
  953. }
  954. /** Map one authoring/roster failure onto its wire code. */
  955. function presetError(agentPreset: string, error: unknown): RpcError {
  956. if (error instanceof UnknownPresetError) {
  957. return {
  958. code: 'agent-preset-not-found',
  959. message: error.message,
  960. details: { agentPreset: error.presetId, available: [...error.available] },
  961. }
  962. }
  963. if (error instanceof PresetNotWritableError) {
  964. return { code: 'agent-preset-read-only', message: error.message, details: { agentPreset, reason: error.message } }
  965. }
  966. if (error instanceof InvalidPresetIdError || error instanceof PresetExistsError) {
  967. return { code: 'agent-preset-invalid', message: error.message, details: { agentPreset, reason: error.message } }
  968. }
  969. return { code: 'internal', message: `agent preset "${agentPreset}": ${String(error)}`, details: {} }
  970. }
  971. class AgentPresetConflict extends Error {
  972. constructor(
  973. readonly sessionId: SessionId,
  974. readonly requestedPreset: string,
  975. readonly existingPreset: string | undefined,
  976. ) {
  977. super(
  978. existingPreset === undefined
  979. ? `session "${sessionId}" records no agent preset, so it cannot be adopted under one; `
  980. + 'a deployment composing no roster records none on any session — '
  981. : `session "${sessionId}" already runs agent preset ${JSON.stringify(existingPreset)}; `
  982. + `requested ${JSON.stringify(requestedPreset)}. A session's preset is fixed at creation.`,
  983. )
  984. }
  985. }
  986. /** Requested identity already belongs to a session with another project cwd. */
  987. class SessionCwdConflict extends Error {
  988. constructor(
  989. readonly sessionId: SessionId,
  990. readonly requestedCwd: string,
  991. readonly existingCwd: string | undefined,
  992. ) {
  993. super(
  994. `session "${sessionId}" already exists with cwd ${JSON.stringify(existingCwd)}; `
  995. + `requested ${JSON.stringify(requestedCwd)}`,
  996. )
  997. }
  998. }
  999. /** An explicit Host naming operation would duplicate another Workspace title. */
  1000. class WorkspaceNameConflictError extends Error {
  1001. constructor(readonly workspaceName: string) {
  1002. super(`workspace name '${workspaceName}' is already in use`)
  1003. this.name = 'WorkspaceNameConflictError'
  1004. }
  1005. }
  1006. /** Shared workspace-not-found error response of the workspace.* mutation rows. */
  1007. function workspaceNotFound<T>(request: RpcRequest<unknown>, workspaceId: string): RpcResponse<T> {
  1008. return err(request, {
  1009. code: 'workspace-not-found',
  1010. message: `workspace "${workspaceId}" not found`,
  1011. details: { workspaceId },
  1012. })
  1013. }
  1014. /** Wire projection of one workspace entity (the workspace.* value row). */
  1015. function workspaceView(workspace: Workspace): WorkspaceView {
  1016. return {
  1017. workspaceId: workspace.id,
  1018. path: workspace.path,
  1019. title: workspace.title,
  1020. sessionIds: [...workspace.sessionIds],
  1021. createdAt: workspace.createdAt,
  1022. updatedAt: workspace.updatedAt,
  1023. }
  1024. }
  1025. /** Wire projection of the durable record carried by `domain/changed`. */
  1026. function changedWorkspaceView(workspaceId: string, value: unknown): WorkspaceView {
  1027. const record: WorkspaceRecord = workspaceRecord.parse(value)
  1028. return {
  1029. workspaceId: workspaceId as WorkspaceId,
  1030. path: record.path,
  1031. title: record.title,
  1032. sessionIds: [...record.sessionIds],
  1033. createdAt: record.createdAt,
  1034. updatedAt: record.updatedAt,
  1035. }
  1036. }
  1037. /** One ApiProxy instance's pending reference-prompt admission listeners. */
  1038. interface PreparedPromptOwnership {
  1039. readonly relocating: Set<MessageId>
  1040. readonly cleanups: Map<MessageId, () => void>
  1041. }
  1042. /** Deliver a prepared prompt and inject its snapshot immediately before that exact message enters. */
  1043. function deliverPrompt(
  1044. ctx: Context,
  1045. agent: Agent,
  1046. mode: 'queue' | 'steer',
  1047. message: UserMessage,
  1048. additionalContext: UserMessage | undefined,
  1049. ownership: PreparedPromptOwnership,
  1050. ): void {
  1051. if (additionalContext === undefined) {
  1052. if (mode === 'steer') agent.steer(message)
  1053. else agent.followup(message)
  1054. return
  1055. }
  1056. let cleanedUp = false
  1057. let detachPreStep = (): void => {}
  1058. let detachDiscard = (): void => {}
  1059. let detachDisposed = (): void => {}
  1060. const cleanup = (): void => {
  1061. /* v8 ignore next -- all settlement paths share this idempotent release. */
  1062. if (cleanedUp) return
  1063. cleanedUp = true
  1064. ownership.cleanups.delete(message.id)
  1065. detachPreStep()
  1066. detachDiscard()
  1067. detachDisposed()
  1068. }
  1069. ownership.cleanups.set(message.id, cleanup)
  1070. // An agent retired with the prepared prompt still pending must not leave
  1071. // these listeners on the Host root context for the process lifetime.
  1072. detachDisposed = ctx.on('agent/disposed', ({ agent: subject }) => {
  1073. if (subject === agent) cleanup()
  1074. })
  1075. detachPreStep = ctx.on('agent/pre-step', async ({ agent: subject, messages }, next): Promise<PreStepDecision> => {
  1076. if (subject !== agent || !messages.some(candidate => candidate.id === message.id)) return next()
  1077. cleanup()
  1078. const decision = await next()
  1079. if (decision.kind !== 'enter') return decision
  1080. const promptIndex = decision.messages.findIndex(candidate => candidate.id === message.id)
  1081. if (promptIndex < 0) return decision
  1082. return {
  1083. kind: 'enter',
  1084. messages: decision.messages.toSpliced(promptIndex, 0, additionalContext),
  1085. }
  1086. }, { prepend: true })
  1087. detachDiscard = ctx.on('agent/inbox/discarded', ({ agent: subject, message: discarded }) => {
  1088. if (subject !== agent || discarded.id !== message.id || ownership.relocating.has(message.id)) return
  1089. const remainsPending = [...agent.inbox.nextTurn, ...agent.inbox.nextStep]
  1090. .some(candidate => candidate.id === message.id)
  1091. if (!remainsPending) cleanup()
  1092. })
  1093. try {
  1094. if (mode === 'steer') agent.steer(message)
  1095. else agent.followup(message)
  1096. } catch (error: unknown) {
  1097. cleanup()
  1098. throw error
  1099. }
  1100. }
  1101. /**
  1102. * Implement ApiProxy over a composed host context.
  1103. * @param ctx - a context with the Host spine and Workspace registry mounted.
  1104. * @param defaults - host routing and project-directory defaults.
  1105. * @returns the ApiProxy implementation.
  1106. */
  1107. export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiProxy {
  1108. const sessionExportCompressionLevel = defaults.sessionExportCompressionLevel
  1109. ?? DEFAULT_SESSION_LOG_COMPRESSION_LEVEL
  1110. const coldBlankProbeMaxBytes = defaults.coldBlankProbeMaxBytes
  1111. ?? DEFAULT_COLD_BLANK_PROBE_MAX_BYTES
  1112. /** The seed model each create/resume declares; re-read so it never goes stale. */
  1113. const agentOptions = (): AgentOptions => {
  1114. const { provider, model } = defaults.defaultModelSelection()
  1115. return { provider, model }
  1116. }
  1117. type WebModelSelectionRef = ModelSelectionRef & { current: ModelSelection }
  1118. const selections = new WeakMap<Agent, WebModelSelectionRef>()
  1119. /**
  1120. * Serializes `agentPreset.select` per session. Two concurrent selects both
  1121. * pass the blank check, and the second `unmountPresetFor` then finds nothing
  1122. * to unmount because the first already removed the record — leaving two
  1123. * compositions registered into one agent layer. The client's `busy` flag is
  1124. * not enforcement: the wire is reachable directly.
  1125. */
  1126. const presetSwitches = new Map<SessionId, Promise<unknown>>()
  1127. /** Client-chosen identity creation/resume, deduplicated across concurrent retries. */
  1128. const sessionCreations = new Map<SessionId, Promise<Agent>>()
  1129. /** Serializes path ownership and explicit title checks with Workspace mutations. */
  1130. let workspaceCreationChain = Promise.resolve()
  1131. const pendingQuestions = new Map<RpcId, PendingQuestion>()
  1132. const pendingApprovals = new Map<RpcId, PendingApproval>()
  1133. const muxQueues = new Set<FrameQueue<RpcRequest<MuxFrame>>>()
  1134. const imageAdmissionChains = new WeakMap<Agent, Promise<void>>()
  1135. const preparedPromptOwnership: PreparedPromptOwnership = {
  1136. relocating: new Set(),
  1137. cleanups: new Map(),
  1138. }
  1139. /** Serialize image admission with model selection for one agent. */
  1140. function serializeImageAdmission<T>(agent: Agent, operation: () => Promise<T>): Promise<T> {
  1141. const result = (imageAdmissionChains.get(agent) ?? Promise.resolve()).then(operation)
  1142. imageAdmissionChains.set(agent, result.then(() => undefined, () => undefined))
  1143. return result
  1144. }
  1145. /**
  1146. * Install or return the session-local model selection that prompt assembly snapshots.
  1147. *
  1148. * Precedence, resolved on EVERY read rather than seeded once: a selection
  1149. * made in this process, else the session's own latest logged request/header,
  1150. * else the live Agent default. Re-reading keeps the two tiers exact in both
  1151. * directions: a session with a recorded request derives its selection from
  1152. * its log, while a blank session (New Session reuses one rather than minting
  1153. * another) reads any default saved after it was created. There is no create-time
  1154. * per-session override tier on this wire — if one returns (a create-options
  1155. * contribution), it must fold in between the selection and the log.
  1156. */
  1157. function selectionFor(agent: Agent): WebModelSelectionRef {
  1158. const installed = selections.get(agent)
  1159. if (installed !== undefined) return installed
  1160. let picked: ModelSelection | undefined
  1161. const selection: WebModelSelectionRef = {
  1162. get current(): ModelSelection {
  1163. if (picked !== undefined) return picked
  1164. // Incrementally folded by the session, so a per-step read costs
  1165. // O(new events) rather than a rescan.
  1166. const logged = agent.session.requestHeader()?.config
  1167. if (logged === undefined) return defaults.defaultModelSelection()
  1168. return {
  1169. provider: logged.provider,
  1170. model: logged.model,
  1171. ...logged.reasoningEffort === undefined
  1172. ? {}
  1173. : { reasoningEffort: logged.reasoningEffort },
  1174. }
  1175. },
  1176. set current(next: ModelSelection) {
  1177. picked = next
  1178. },
  1179. assembled: undefined,
  1180. }
  1181. installModelSelection(agent.ctx, selection)
  1182. selections.set(agent, selection)
  1183. return selection
  1184. }
  1185. /** Pre-publication setup used by both fresh and resumed Web agents. */
  1186. function installSelection(agentCtx: Context): void {
  1187. const agent = agentCtx.agent
  1188. if (agent === undefined) throw new Error('api-proxy: agent setup has no scoped agent')
  1189. selectionFor(agent)
  1190. }
  1191. /**
  1192. * Reject an attempt to run an existing session under a different preset.
  1193. *
  1194. * A caller that names no preset always adopts the session as it is, so the
  1195. * common paths — reconnecting, resuming, retrying a create — are unaffected.
  1196. * @param sessionId - the identity being adopted.
  1197. * @param requested - the preset the request named, if any.
  1198. * @param existing - the preset the session RUNS, if any; both callers resolve
  1199. * it from the log, which differs from the creation header once a blank
  1200. * session has switched.
  1201. * @throws when both are present and differ.
  1202. */
  1203. function assertPresetUnchanged(
  1204. sessionId: SessionId,
  1205. requested: string | undefined,
  1206. existing: string | undefined,
  1207. ): void {
  1208. if (requested === undefined || requested === existing) return
  1209. throw new AgentPresetConflict(sessionId, requested, existing)
  1210. }
  1211. /**
  1212. * Resolve the preset an agent will be composed from, and the setup that
  1213. * installs it.
  1214. *
  1215. * The id is resolved BEFORE the session exists because the session boundary
  1216. * snapshots `meta` before asynchronous setup begins — a preset discovered
  1217. * during setup could never reach the header. Mounting still happens in
  1218. * setup, where a failure rolls the whole creation back rather than leaving a
  1219. * published session whose capabilities are half-installed.
  1220. *
  1221. * A deployment with no preset roster composes nothing and every session
  1222. * shares the host composition, which is the behavior before presets existed.
  1223. * @param presetId - the requested preset, or `undefined` for the default.
  1224. * @returns the id to record on the header (absent without a roster) and the setup callback.
  1225. * @throws when the roster supplies no such preset.
  1226. */
  1227. async function composeAgent(presetId: string | undefined): Promise<{
  1228. agentPreset?: string
  1229. setup: (agentCtx: Context) => Promise<void>
  1230. }> {
  1231. const presets = ctx.get('agentPresets')
  1232. if (presets === undefined) {
  1233. return {
  1234. setup: (agentCtx: Context) => {
  1235. installSelection(agentCtx)
  1236. return Promise.resolve()
  1237. },
  1238. }
  1239. }
  1240. const resolvedId = (await presets.resolve(presetId)).id
  1241. return {
  1242. agentPreset: resolvedId,
  1243. setup: async (agentCtx: Context) => {
  1244. installSelection(agentCtx)
  1245. await presets.mount(agentCtx, resolvedId)
  1246. },
  1247. }
  1248. }
  1249. const hasSubagentOwner = (
  1250. session: Pick<Session, 'header'>,
  1251. agent: Agent | undefined,
  1252. ): boolean => hasApiRemoteSubagentOwner(ctx, session, agent)
  1253. const subagentOwnershipError = (sessionId: SessionId): RpcError =>
  1254. apiRemoteSubagentOwnershipError(sessionId)
  1255. const inspectServable = (sessionId: SessionId): Promise<{ meta: SessionHeader; events: SessionEvent[] }> =>
  1256. inspectApiRemoteSession(ctx, sessionId)
  1257. // Cold resume composes the preset the session recorded, for the same reason
  1258. // `session.create` does: its history was produced under that composition.
  1259. // Every generic entry point — prompt, models, commands — arrives here, so
  1260. // leaving it out meant a session opened after a restart ran on host tools
  1261. // and the deployment persona. Resolved from the LOG, not the header: a
  1262. // session that switched while blank ran its turns under the newer
  1263. // composition, and the header is written once at creation. Reading the
  1264. // header here would silently undo the switch on the next restart and
  1265. // restore that history under the old tool set.
  1266. const agentFor = createApiRemoteAgentResolver(ctx, {
  1267. agentOptions,
  1268. setup: async ({ meta, events }) =>
  1269. (await composeAgent(resolveSessionPreset({ header: meta, events }))).setup,
  1270. })
  1271. /** Send one transient frame to every connected mux consumer. */
  1272. function broadcast(payload: MuxFrame): void {
  1273. const envelope = frame(payload)
  1274. for (const queue of muxQueues) queue.push(envelope)
  1275. }
  1276. // Projection change feed → session/projection push frames. The carrier
  1277. // mints the wire frame (the Service Definition package holds no wire vocabulary); the
  1278. // child activates only when a projection registry is composed, and the
  1279. // subscription unwinds with this gateway's fiber.
  1280. ctx.inject(['sessionProjections'], (projectionCtx) => {
  1281. projectionCtx.sessionProjections.onChanged((session, key, value, seq) => {
  1282. broadcast({ type: 'session/projection', sessionId: session.id, key, value, seq })
  1283. })
  1284. })
  1285. // The cache supplies recency and a monotonic non-blank hint. A cached
  1286. // `blank: true` remains only a prefix fact and is verified on the cold path.
  1287. ctx.inject(['sessionProjections'], (projectionCtx) => {
  1288. projectionCtx.sessionProjections.register<'sessionListMetadata', SessionListMetadata>({
  1289. key: 'sessionListMetadata',
  1290. schema: sessionListMetadataProjectionSchema,
  1291. init: () => ({ blank: true, lastPromptAt: null }),
  1292. apply: applySessionListMetadata,
  1293. view: state => state,
  1294. stateVersion: 1,
  1295. })
  1296. })
  1297. // The imageLimits projection unit: the attachments config this proxy
  1298. // enforces at prompt admission, constant per host boot. `apply` keeps the
  1299. // same state reference for every event, so no change frames are ever
  1300. // pushed — baselines alone carry the value — and clients pre-check intake
  1301. // and label upload affordances from it. Registered here, not in the
  1302. // attachment Service Definition: dsh-llm depends on dsh-attachment, so the
  1303. // seam package cannot reference the projection registry without a cycle,
  1304. // and the per-message rules the value describes are this proxy's own
  1305. // admission checks. The child activates only while both seams are composed.
  1306. // `view` reading the live service instead of the (null) state is sanctioned
  1307. // exactly for boot-constant units: the value cannot change within a process
  1308. // lifetime, so the fold stays observationally pure, and a stale persisted
  1309. // cache row re-viewing to the current config is the correct outcome.
  1310. ctx.inject(['sessionProjections', 'attachments'], (projectionCtx) => {
  1311. projectionCtx.sessionProjections.register<'imageLimits', null>({
  1312. key: 'imageLimits',
  1313. schema: imageLimitsProjectionSchema,
  1314. init: () => null,
  1315. apply: state => state,
  1316. view: () => projectionCtx.attachments.imageLimits,
  1317. stateVersion: 1,
  1318. })
  1319. })
  1320. /** Project both durable inbox lists, optionally including the splice currently being emitted. */
  1321. const queueItems = (
  1322. agent: Agent,
  1323. splice?: SessionEventMap['agent/inbox/spliced'],
  1324. ): QueuedInboxItem[] => {
  1325. const project = (target: 'next-turn' | 'next-step'): readonly UserMessage[] => {
  1326. const messages = target === 'next-turn' ? agent.inbox.nextTurn : agent.inbox.nextStep
  1327. return splice?.target === target
  1328. ? messages.toSpliced(splice.start, splice.removedCount ?? 0, ...splice.inserted)
  1329. : messages
  1330. }
  1331. return [
  1332. ...project('next-turn').map(message => ({ id: message.id, placement: 'queued' as const, message })),
  1333. ...project('next-step').map(message => ({
  1334. id: message.id,
  1335. // Only user-origin messages are steering; injected context (approval
  1336. // notices, task completion, attached snapshots) is not a user action
  1337. // and must not render as a pending steering bubble.
  1338. placement: message.source.kind === 'user' ? 'steering' as const : 'context' as const,
  1339. message,
  1340. })),
  1341. ]
  1342. }
  1343. ctx.on('session/event', (session, event) => {
  1344. if (event.type !== 'agent/inbox/spliced') return
  1345. const agent = ctx.agents.get(session.id)
  1346. if (agent?.session !== session) return
  1347. broadcast({ type: 'session/queue', sessionId: session.id, items: queueItems(agent, event.data) })
  1348. })
  1349. /** Remove a wait before settling it: synchronous deletion makes the first claimant win. */
  1350. function claimQuestion(pending: PendingQuestion, outcome: 'answered' | 'cancelled'): void {
  1351. pendingQuestions.delete(pending.rpcId)
  1352. if (pending.signal !== undefined && pending.onAbort !== undefined) {
  1353. pending.signal.removeEventListener('abort', pending.onAbort)
  1354. }
  1355. broadcast({
  1356. type: 'question/resolved', sessionId: pending.sessionId,
  1357. questionRpcId: pending.rpcId, outcome,
  1358. })
  1359. }
  1360. const disposeProvider = ctx.userQuestions.registerProvider({
  1361. ask(request: AskUserQuestionRequest): Promise<AskUserQuestionAnswer> {
  1362. const sessionId = request.agent?.id
  1363. if (sessionId === undefined) {
  1364. return Promise.reject(new UserQuestionError(
  1365. 'web user interaction requires an agent-owned session', 'ASK_MISSING_AGENT'))
  1366. }
  1367. return new Promise<AskUserQuestionAnswer>((resolve, reject) => {
  1368. const rpcId = RpcId(randomUUID())
  1369. const pending: PendingQuestion = {
  1370. rpcId, sessionId, questions: request.questions, resolve, reject,
  1371. ...(request.signal === undefined ? {} : { signal: request.signal }),
  1372. }
  1373. const onAbort = (): void => {
  1374. claimQuestion(pending, 'cancelled')
  1375. reject(new UserQuestionError(
  1376. 'ask_user_question was aborted before the user answered', 'ASK_ABORTED'))
  1377. }
  1378. pending.onAbort = onAbort
  1379. pendingQuestions.set(rpcId, pending)
  1380. request.signal?.addEventListener('abort', onAbort, { once: true })
  1381. const envelope: RpcRequest<MuxFrame> = {
  1382. rpcId,
  1383. payload: { type: 'question/requested', sessionId, questions: request.questions },
  1384. }
  1385. for (const queue of muxQueues) queue.push(envelope)
  1386. })
  1387. },
  1388. })
  1389. ctx.effect(() => () => {
  1390. disposeProvider()
  1391. for (const pending of [...pendingQuestions.values()]) {
  1392. claimQuestion(pending, 'cancelled')
  1393. pending.reject(new UserQuestionError(
  1394. 'web user-questions provider was disposed', 'ASK_ABORTED'))
  1395. }
  1396. }, 'api-proxy: user-questions provider')
  1397. // --- Approval pending registry ------------------------------------------
  1398. // The proxy is the approval channel for every agent this host owns: an ask
  1399. // through `ctx.approval` becomes an answerable server-request on the mux
  1400. // stream (stable rpcId), settled by POST /api/respond. The entry survives
  1401. // client disconnects — mux-open replays still-pending requested frames with
  1402. // the same rpcId (the refresh-recovery baseline) — and withdraws on the
  1403. // ask's own abort signal (turn cancel), pushing `cancelled` to subscribers.
  1404. if (ctx.get('approval') !== undefined) {
  1405. // Teardown parity with the question provider above: a gateway disposed
  1406. // while approvals are pending settles every entry as 'cancelled' (the
  1407. // service's fail-closed vocabulary), so no ask promise dangles past the
  1408. // proxy's lifetime and subscribers see the withdrawal.
  1409. ctx.effect(() => () => {
  1410. for (const pending of [...pendingApprovals.values()]) pending.resolve('cancelled')
  1411. }, 'api-proxy: approval registry teardown')
  1412. ctx.on('approval/request', (req, next) => {
  1413. // Dispatch rides a microtask behind the service's own signal check: an
  1414. // abort landing in that window would register the abort listener AFTER
  1415. // the signal fired — never invoked, entry pending forever, zombie frame
  1416. // on every mux replay. Settle synchronously instead of publishing.
  1417. if (req.signal?.aborted === true) return Promise.resolve<ApprovalOutcome>('cancelled')
  1418. // The audit pair `approval/asked` is already appended by the service
  1419. // before dispatch, but dispatch rides a microtask: parallel tool calls
  1420. // can append several asked events before any answerer runs. THIS
  1421. // request's event is therefore the newest asked event that is still
  1422. // undecided, unclaimed by another pending entry, and — when the ask
  1423. // names a call — carries the same callId.
  1424. const events = req.agent.session.events
  1425. const claimed = new Set<ApprovalRequestId>()
  1426. for (const entry of pendingApprovals.values()) claimed.add(entry.approvalId)
  1427. const decided = new Set<ApprovalRequestId>()
  1428. let approvalId: ApprovalRequestId | undefined
  1429. for (let i = events.length - 1; i >= 0; i -= 1) {
  1430. const event = events[i] as SessionEvent
  1431. if (event.type === 'approval/decided') {
  1432. decided.add(event.data.id)
  1433. } else if (event.type === 'approval/asked') {
  1434. if (decided.has(event.data.id) || claimed.has(event.data.id)) continue
  1435. // Symmetric pairing: a callId-bearing ask only takes its own call's
  1436. // record, and a callId-less ask only takes a callId-less record —
  1437. // so neither shape can steal the other's audit id under parallel
  1438. // asks. (Today every producer — the tool executor — passes callId;
  1439. // the callId-less arm guards any future non-tool asker.)
  1440. if ((req.callId ?? null) !== (event.data.callId ?? null)) continue
  1441. approvalId = event.data.id
  1442. break
  1443. }
  1444. }
  1445. // No asked event means the request bypassed the service's audit path —
  1446. // not this channel's question; delegate to the fail-closed default.
  1447. if (approvalId === undefined) return next()
  1448. const id = approvalId
  1449. return new Promise<ApprovalOutcome>((resolve) => {
  1450. const settle = (outcome: ApprovalOutcome): void => {
  1451. /* v8 ignore next 3 -- defensive double-settle guard: respond() routes
  1452. through the pending table (a settled id is not-pending before it can
  1453. re-settle) and the first settle removes the abort listener, so no
  1454. reachable path settles twice; kept against future settle callers. */
  1455. if (!pendingApprovals.delete(pending.rpcId)) return
  1456. req.signal?.removeEventListener('abort', onAbort)
  1457. broadcast({ type: 'approval/resolved', sessionId: pending.sessionId, approvalId: id, outcome })
  1458. // A cancelled ask was already settled by the service's own signal
  1459. // race, which discards this late resolution; resolving is a no-op
  1460. // there and keeps this promise from dangling forever.
  1461. resolve(outcome)
  1462. }
  1463. const onAbort = (): void => { settle('cancelled') }
  1464. const pending: PendingApproval = {
  1465. rpcId: RpcId(randomUUID()),
  1466. sessionId: req.agent.session.id,
  1467. approvalId: id,
  1468. toolName: req.toolName,
  1469. ...req.callId === undefined ? {} : { callId: req.callId },
  1470. ...req.reason === undefined ? {} : { reason: req.reason },
  1471. resolve: settle,
  1472. }
  1473. pendingApprovals.set(pending.rpcId, pending)
  1474. req.signal?.addEventListener('abort', onAbort, { once: true })
  1475. const envelope = requestedFrame(pending)
  1476. for (const queue of muxQueues) queue.push(envelope)
  1477. })
  1478. })
  1479. }
  1480. type SessionReadState = {
  1481. id: SessionId
  1482. header: SessionHeader
  1483. events: SessionEvent[]
  1484. }
  1485. /** Read one stable session prefix without acquiring an Agent owner. */
  1486. async function readSessionState(sessionId: SessionId): Promise<SessionReadState> {
  1487. const attached = ctx.sessions.get(sessionId)
  1488. if (attached !== undefined) {
  1489. return {
  1490. id: attached.id,
  1491. header: attached.header,
  1492. events: [...attached.events],
  1493. }
  1494. }
  1495. const inspected = await inspectServable(sessionId)
  1496. return { id: inspected.meta.id, header: inspected.meta, events: inspected.events }
  1497. }
  1498. /** Resolve the Workspace inherited by a fork without making ordinary loose lineage grouped. */
  1499. async function forkWorkspace(source: Pick<Session, 'id' | 'header'>): Promise<Workspace | undefined> {
  1500. const workspaces = ctx.workspaceRegistry.list()
  1501. const direct = workspaces.find(workspace => workspace.sessionIds.includes(source.id))
  1502. if (direct !== undefined || source.header.origin !== 'subagent') return direct
  1503. const lineage = await ctx.sessionQuery.traceSession(source.id)
  1504. for (const ancestor of lineage.ancestors) {
  1505. const workspace = workspaces.find(candidate => candidate.sessionIds.includes(ancestor.header.id))
  1506. if (workspace !== undefined) return workspace
  1507. }
  1508. return undefined
  1509. }
  1510. /**
  1511. * Resolve which session one transcript read is served from, without
  1512. * acquiring an Agent owner. This is the read's only asynchronous step
  1513. * besides ensuring the composition; {@link historyCutOf} takes the cut.
  1514. * @param sessionId - the transcript being read.
  1515. * @returns the attached session, or the inspected detached header and events.
  1516. * @throws {@link ApiRemoteSessionNotFound} when no project-backed session has that identity.
  1517. */
  1518. async function historySourceFor(sessionId: SessionId): Promise<HistorySource> {
  1519. const attached = ctx.sessions.get(sessionId)
  1520. if (attached !== undefined) return { kind: 'attached', session: attached }
  1521. const inspected = await inspectServable(sessionId)
  1522. return { kind: 'detached', header: inspected.meta, events: inspected.events }
  1523. }
  1524. /**
  1525. * The header and events {@link presenterScopeFor} reads to decide which
  1526. * composition a transcript ran under.
  1527. * @param source - the live or detached session this read is served from.
  1528. * @returns that session's creation header and its events.
  1529. */
  1530. function sourceSession(source: HistorySource): PresetBearingSession {
  1531. if (source.kind === 'detached') return { header: source.header, events: source.events }
  1532. return { header: source.session.header, events: source.session.events }
  1533. }
  1534. /**
  1535. * One transcript cut: the events and the projection baseline that describe
  1536. * the SAME log position.
  1537. *
  1538. * Synchronous, and the two reads sit next to each other, because an attached
  1539. * session keeps appending: an `await` between them would serve events cut at
  1540. * N beside a baseline folded to N+1, which is one response describing two
  1541. * moments. The caller does its awaiting before this call.
  1542. * @param source - the live or detached session this read is served from.
  1543. * @param includeProjections - whether the caller asked for the baseline (a tail page does).
  1544. * @returns the events and, when asked, the baseline for that same position.
  1545. */
  1546. function historyCutOf(
  1547. source: HistorySource,
  1548. includeProjections: boolean,
  1549. ): { events: SessionEvent[]; projections?: SessionProjectionsBlock } {
  1550. if (source.kind === 'detached') {
  1551. const projections = includeProjections ? detachedProjectionsFor(ctx, source.events) : undefined
  1552. return { events: source.events, ...projections === undefined ? {} : { projections } }
  1553. }
  1554. const events = [...source.session.events]
  1555. const projections = includeProjections ? projectionsFor(ctx, source.session) : undefined
  1556. return { events, ...projections === undefined ? {} : { projections } }
  1557. }
  1558. /**
  1559. * The registry view scope a transcript's presenters resolve in.
  1560. *
  1561. * A live agent is that scope itself (its chain passes through its preset's
  1562. * standing layer). A cold session resolves its preset from the LOG, and the
  1563. * preset's STANDING key serves without resuming anything — ensuring the
  1564. * mount composes plugins but starts no agent, session, or turn. No roster,
  1565. * no recorded preset, or a preset the roster no longer supplies all fall
  1566. * back to the global layer: the transcript still serves, with the generic
  1567. * cards a viewless entry renders.
  1568. *
  1569. * Reading the header alone would render a session that switched while blank
  1570. * through the composition it was CREATED with. Every tool only the newer
  1571. * preset registers resolves to no presenter there, and the transcript
  1572. * silently degrades to generic cards for exactly the calls its history is
  1573. * made of.
  1574. * @param sessionId - the transcript being read.
  1575. * @param session - that session's header and log (attached or inspected).
  1576. * @returns the scope to pass to presenter lookups, or undefined for global.
  1577. */
  1578. async function presenterScopeFor(
  1579. sessionId: SessionId,
  1580. session: PresetBearingSession,
  1581. ): Promise<ScopeKey | undefined> {
  1582. const live = ctx.get('agents')?.get(sessionId)
  1583. if (live !== undefined) return live
  1584. const presets = ctx.get('agentPresets')
  1585. if (presets === undefined) return undefined
  1586. try {
  1587. // An unrecorded preset (a log from before the roster existed) renders
  1588. // through the DEFAULT preset's standing layer: that is the composition
  1589. // an unnamed session composes today, and presenters are pure display,
  1590. // so the worst a mismatch produces is the generic card it had anyway.
  1591. return await presets.standingKeyFor(resolveSessionPreset(session))
  1592. } catch {
  1593. // Swallows only the unknown/unusable-preset rejection from the roster:
  1594. // a deleted or broken preset must degrade this read, never fail it.
  1595. return undefined
  1596. }
  1597. }
  1598. /** Resolve one requested identity to a live agent, creating or resuming it once. */
  1599. async function ensureSession(
  1600. sessionId: SessionId,
  1601. cwd: string,
  1602. checkPersistedIdentity: boolean,
  1603. presetId?: string,
  1604. ): Promise<Agent> {
  1605. let creation = sessionCreations.get(sessionId)
  1606. if (creation === undefined) {
  1607. creation = (async () => {
  1608. const attached = ctx.sessions.get(sessionId)
  1609. const live = ctx.agents.get(sessionId)
  1610. if (attached !== undefined && hasSubagentOwner(attached, live)) {
  1611. throw new SubagentSessionOwnership(sessionId)
  1612. }
  1613. if (live !== undefined) return live
  1614. const persistence = checkPersistedIdentity ? ctx.get('sessionPersistence') : undefined
  1615. const stored = persistence === undefined
  1616. ? undefined
  1617. : (await persistence.list()).find(header => header.id === sessionId)
  1618. if (persistence !== undefined && stored !== undefined) {
  1619. const inspected = await persistence.inspect(sessionId)
  1620. // Ownership first: explicit-id adoption of a session-backed
  1621. // subagent must answer `agent-busy` regardless of the requested
  1622. // cwd (the api/commands.ts contract), not a cwd conflict.
  1623. if (hasSubagentOwner({ header: inspected.meta }, undefined)) {
  1624. throw new SubagentSessionOwnership(sessionId)
  1625. }
  1626. if (inspected.meta.cwd !== cwd) {
  1627. throw new SessionCwdConflict(sessionId, cwd, inspected.meta.cwd)
  1628. }
  1629. // Resolved from the log, not the header: a session that switched
  1630. // while blank ran every turn under the newer composition.
  1631. const storedPreset = resolveSessionPreset({ header: inspected.meta, events: inspected.events })
  1632. assertPresetUnchanged(sessionId, presetId, storedPreset)
  1633. // The stored preset wins over anything the request names: a resumed
  1634. // session's history was produced under that composition, and
  1635. // rebuilding it differently would replay tool calls the model can no
  1636. // longer make.
  1637. return (await ctx.agents.resume({
  1638. resumeSessionId: sessionId,
  1639. agentOptions: agentOptions(),
  1640. setup: (await composeAgent(storedPreset)).setup,
  1641. })).agent
  1642. }
  1643. try {
  1644. await mkdir(cwd, { recursive: true })
  1645. } catch (error: unknown) {
  1646. throw new Error(`failed to ensure project directory "${cwd}": ${String(error)}`, { cause: error })
  1647. }
  1648. const composition = await composeAgent(presetId)
  1649. return (await ctx.agents.create({
  1650. sessionId,
  1651. agentOptions: agentOptions(),
  1652. meta: {
  1653. cwd,
  1654. ...composition.agentPreset === undefined ? {} : { agentPreset: composition.agentPreset },
  1655. },
  1656. setup: composition.setup,
  1657. })).agent
  1658. })().catch((error: unknown) => {
  1659. // Another Host entry path may have published the same identity while
  1660. // this operation crossed an asynchronous persistence/filesystem step.
  1661. const live = ctx.agents.get(sessionId)
  1662. if (live !== undefined) {
  1663. if (hasSubagentOwner(live.session, live)) throw new SubagentSessionOwnership(sessionId)
  1664. return live
  1665. }
  1666. const attached = ctx.sessions.get(sessionId)
  1667. if (attached !== undefined && hasSubagentOwner(attached, undefined)) {
  1668. throw new SubagentSessionOwnership(sessionId)
  1669. }
  1670. throw error
  1671. }).finally(() => {
  1672. sessionCreations.delete(sessionId)
  1673. })
  1674. sessionCreations.set(sessionId, creation)
  1675. }
  1676. const agent = await creation
  1677. if (hasSubagentOwner(agent.session, agent)) throw new SubagentSessionOwnership(sessionId)
  1678. // Beside the cwd check for the same reason, and after the await so it
  1679. // covers every path that yields a live agent — freshly created, adopted
  1680. // live, resumed from disk, or recovered by the concurrent-creation catch.
  1681. assertPresetUnchanged(sessionId, presetId, resolveSessionPreset(agent.session))
  1682. if (agent.session.header.cwd !== cwd) {
  1683. throw new SessionCwdConflict(sessionId, cwd, agent.session.header.cwd)
  1684. }
  1685. return agent
  1686. }
  1687. /** Resolve or create one path while holding the Host's workspace-create chain. */
  1688. function ensureWorkspace(path: string): Promise<{ workspace: Workspace; created: boolean }> {
  1689. const operation = workspaceCreationChain.then(async () => {
  1690. const existing = await ctx.workspaceRegistry.resolveByPath(path)
  1691. if (existing !== undefined) return { workspace: existing, created: false }
  1692. return { workspace: await ctx.workspaceRegistry.create(path), created: true }
  1693. })
  1694. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  1695. return operation
  1696. }
  1697. /**
  1698. * Build the session.list baseline shared by listing and search visibility.
  1699. * Attached sessions come from memory; servable cold sessions merge from
  1700. * persistence, and the final order is newest-first.
  1701. */
  1702. async function listVisibleSessionSummaries(signal?: AbortSignal): Promise<SessionSummary[]> {
  1703. signal?.throwIfAborted()
  1704. const summarizeAttached = (session: Session): SessionSummary => {
  1705. const agent = ctx.agents.get(session.id)
  1706. const projections = listProjectionsFor(ctx, session.header, session)
  1707. return {
  1708. ...summarize(session, agent?.status === 'running'),
  1709. ...projections === undefined ? {} : { projections },
  1710. }
  1711. }
  1712. const items = ctx.sessions.list().map(summarizeAttached)
  1713. signal?.throwIfAborted()
  1714. const attached = new Set(items.map(item => item.sessionId))
  1715. const persistence = ctx.get('sessionPersistence')
  1716. if (persistence !== undefined) {
  1717. const cold = (await persistence.list(signal))
  1718. .filter(meta => !attached.has(meta.id) && meta.cwd !== undefined)
  1719. signal?.throwIfAborted()
  1720. for (let offset = 0; offset < cold.length; offset += COLD_SUMMARY_BATCH_SIZE) {
  1721. signal?.throwIfAborted()
  1722. const batch = cold.slice(offset, offset + COLD_SUMMARY_BATCH_SIZE)
  1723. const settled = await Promise.allSettled(
  1724. batch.map(async (meta) => {
  1725. // Projection hints remain optional. Blank verification may read
  1726. // this Session's artifact only when it passes the configured size check.
  1727. const projections = listProjectionsFor(ctx, meta, undefined)
  1728. const summary = await summarizeCold(
  1729. ctx,
  1730. persistence,
  1731. meta,
  1732. projections?.values.sessionListMetadata,
  1733. coldBlankProbeMaxBytes,
  1734. signal,
  1735. )
  1736. const attachedSession = ctx.sessions.get(meta.id)
  1737. if (attachedSession !== undefined) return summarizeAttached(attachedSession)
  1738. return {
  1739. ...summary,
  1740. ...projections === undefined ? {} : { projections },
  1741. }
  1742. }),
  1743. )
  1744. const summaries: SessionSummary[] = []
  1745. let rejected = false
  1746. let failure: unknown
  1747. for (const result of settled) {
  1748. if (result.status === 'fulfilled') {
  1749. summaries.push(result.value)
  1750. } else if (!rejected) {
  1751. rejected = true
  1752. failure = result.reason
  1753. }
  1754. }
  1755. if (rejected) throw failure
  1756. signal?.throwIfAborted()
  1757. items.push(...summaries)
  1758. }
  1759. }
  1760. items.sort((a, b) => b.updatedAt - a.updatedAt)
  1761. return items
  1762. }
  1763. /**
  1764. * Resolve the goal service THIS agent runs.
  1765. *
  1766. * The service is per session: an agent preset mounts it behind an `isolate`
  1767. * realm, which no host context resolves. Reading it from the root would
  1768. * answer "absent" for a session whose composition mounts it — so the lookup
  1769. * is keyed by the agent, and only a deployment composing it nowhere is
  1770. * genuinely absent.
  1771. */
  1772. function goalServiceFor(agent: Agent): NonNullable<ReturnType<typeof ctx.get<'goals'>>> | { error: RpcError } {
  1773. const presets = ctx.get('agentPresets')
  1774. const goals = presets?.serviceFor(agent, 'goals') ?? ctx.get('goals')
  1775. if (goals === undefined) {
  1776. return { error: { code: 'internal', message: 'goal service is absent: neither this session\'s agent preset nor the host composition mounts @deepseek-ai/dsh-goal', details: {} } }
  1777. }
  1778. return goals
  1779. }
  1780. /** Map one goal-domain rejection to the wire error (stable GoalError codes ride in details). */
  1781. function goalError(request: RpcRequest<unknown>, error: unknown): RpcResponse<never> {
  1782. const details = error instanceof GoalError ? { goalCode: error.code } : {}
  1783. return err(request, { code: 'internal', message: String(error), details })
  1784. }
  1785. /** Resolve a session's agent, apply one goal mutation, and acknowledge with the new CAS ref. */
  1786. async function mutateGoal(
  1787. request: RpcRequest<{ sessionId: SessionId }>,
  1788. mutation: (goals: NonNullable<ReturnType<typeof ctx.get<'goals'>>>, agent: Agent) => CoreGoalRef,
  1789. ): Promise<RpcResponse<{ ref: GoalRef }>> {
  1790. const found = await agentFor(request.payload.sessionId)
  1791. if ('error' in found) return err(request, found.error)
  1792. const goals = goalServiceFor(found.agent)
  1793. if ('error' in goals) return err(request, goals.error)
  1794. try {
  1795. const ref = mutation(goals, found.agent)
  1796. return ok(request, { ref: { id: ref.id, revision: ref.revision } })
  1797. } catch (error: unknown) {
  1798. return goalError(request, error)
  1799. }
  1800. }
  1801. /**
  1802. * Whether an adapter currently serves this provider, and therefore whether
  1803. * a session selecting it can start a turn. Catalog membership cannot answer
  1804. * it: an adapter may serve a model its own catalog stopped advertising, so
  1805. * a provider missing from the groups is not the same as one nothing serves.
  1806. * A composition with no llm registry at all cannot judge and says yes —
  1807. * the dispatch it would have refused fails on its own terms.
  1808. */
  1809. function routeServed(provider: string): boolean {
  1810. const llm = ctx.get('llm')
  1811. return llm === undefined || llm.listProviders().some(entry => entry.id === provider)
  1812. }
  1813. /**
  1814. * Resolve the addressed agent for a turn-starting method and refuse when no
  1815. * adapter serves its current selection: a provider nothing serves cannot start a
  1816. * turn, and letting it try spends the whole pre-step path to fail inside
  1817. * the adapter with a message about registration. Refusing here names the
  1818. * model the session is pointed at while the draft is still in the composer.
  1819. * This is `session.prompt`'s enforcement boundary: a client that disables
  1820. * its input is an affordance, and the method stays callable regardless.
  1821. */
  1822. async function turnAgentFor<T>(
  1823. request: RpcRequest<unknown>, sessionId: SessionId,
  1824. ): Promise<{ agent: Agent } | { refused: RpcResponse<T> }> {
  1825. const found = await agentFor(sessionId)
  1826. if ('error' in found) return { refused: err(request, found.error) }
  1827. const agent = found.agent
  1828. const selection = selectionFor(agent).current
  1829. if (!routeServed(selection.provider)) {
  1830. return {
  1831. refused: err(request, {
  1832. code: 'model-unavailable',
  1833. message: `no adapter serves provider "${selection.provider}"; select a model for this session`,
  1834. details: { provider: selection.provider, model: selection.model },
  1835. }),
  1836. }
  1837. }
  1838. return { agent }
  1839. }
  1840. /** Missing-service report shared by the settings domain (skills-domain stance). */
  1841. function settingsAbsent(): RpcError {
  1842. return { code: 'internal', message: 'settings service is absent: this deployment does not mount a settings provider (e.g. @deepseek-ai/dsh-settings-file) in its composition', details: {} }
  1843. }
  1844. /** Open one Host-resolved target and map native failures onto the wire vocabulary. */
  1845. async function openTarget(
  1846. request: RpcRequest<unknown>, path: string, signal: AbortSignal,
  1847. open: (path: string, signal: AbortSignal) => Promise<void>,
  1848. ): Promise<RpcResponse<{ opened: true }>> {
  1849. try {
  1850. await open(path, signal)
  1851. return ok(request, { opened: true as const })
  1852. } catch (error: unknown) {
  1853. if (signal.aborted) {
  1854. return err(request, {
  1855. code: 'cancelled',
  1856. message: 'path open was aborted',
  1857. details: {},
  1858. })
  1859. }
  1860. return err(request, {
  1861. code: 'internal',
  1862. message: `path open failed: ${error instanceof Error ? error.message : String(error)}`,
  1863. details: {},
  1864. })
  1865. }
  1866. }
  1867. /** Open one Host-resolved path with its default application. */
  1868. function openPath(
  1869. request: RpcRequest<unknown>, path: string, signal: AbortSignal,
  1870. ): Promise<RpcResponse<{ opened: true }>> {
  1871. const open = defaults.openPath
  1872. ?? ((target: string, openSignal: AbortSignal) => openNativePath(target, openSignal))
  1873. return openTarget(request, path, signal, open)
  1874. }
  1875. /** Open one Host-resolved text document in a native editor. */
  1876. function openTextFile(
  1877. request: RpcRequest<unknown>, path: string, signal: AbortSignal,
  1878. ): Promise<RpcResponse<{ opened: true }>> {
  1879. const open = defaults.openTextFile
  1880. ?? ((target: string, openSignal: AbortSignal) => openNativeTextFile(target, openSignal))
  1881. return openTarget(request, path, signal, open)
  1882. }
  1883. /** Whether this deployment can hand a path to a native opener at all. */
  1884. function canOpenPaths(): boolean {
  1885. if (defaults.canOpenPath !== undefined) return defaults.canOpenPath()
  1886. // An injected opener is by definition usable; otherwise ask the platform.
  1887. return defaults.openPath !== undefined || canOpenNativePath()
  1888. }
  1889. /** Missing-service report shared by the credentials domain. */
  1890. function credentialsAbsent(): RpcError {
  1891. 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: {} }
  1892. }
  1893. /** Map one redacted settings descriptor to its wire view. */
  1894. function namespaceView(descriptor: SettingsDescriptor): SettingsNamespaceView {
  1895. return {
  1896. ns: String(descriptor.ns),
  1897. schema: descriptor.schema,
  1898. value: descriptor.value,
  1899. ...descriptor.base === undefined ? {} : { base: descriptor.base },
  1900. ...descriptor.user === undefined ? {} : { user: descriptor.user },
  1901. applies: descriptor.applies,
  1902. secrets: (descriptor.secrets ?? []).map(secret => ({ path: [...secret.path], set: secret.set })),
  1903. revision: descriptor.revision,
  1904. }
  1905. }
  1906. /**
  1907. * Run one settings write (merge or wholesale replace) and acknowledge with
  1908. * the namespace's new redacted view. Every seam refusal — unknown or invalid
  1909. * namespace, read-only provider, schema validation, storage — becomes one
  1910. * `settings-rejected` carrying the seam's own message.
  1911. */
  1912. async function settingsWrite(
  1913. request: RpcRequest<unknown>,
  1914. ns: string,
  1915. mode: 'update' | 'replace' | 'mutate',
  1916. section: object,
  1917. expectedRevision?: number,
  1918. ): Promise<RpcResponse<SettingsNamespaceView>> {
  1919. const settings = ctx.get('settings')
  1920. if (settings === undefined) return err(request, settingsAbsent())
  1921. const rejected = (error: unknown): RpcResponse<SettingsNamespaceView> => {
  1922. // A stale writer is its own outcome, not a malformed request: the client
  1923. // must re-read and re-apply rather than treat the write as invalid.
  1924. if (error instanceof SettingsConflictError) {
  1925. return err(request, {
  1926. code: 'settings-conflict',
  1927. message: error.message,
  1928. details: { ns, expected: error.expected, actual: error.actual },
  1929. })
  1930. }
  1931. return err(request, {
  1932. code: 'settings-rejected',
  1933. message: error instanceof Error ? error.message : String(error),
  1934. details: { ns },
  1935. })
  1936. }
  1937. let branded: SettingsNamespace
  1938. try {
  1939. branded = settingsNamespace(ns)
  1940. } catch (error: unknown) {
  1941. // A malformed name can address no registration, so it fails exactly as
  1942. // an unregistered one does.
  1943. return rejected(error)
  1944. }
  1945. try {
  1946. if (mode === 'update') await settings.update(branded, section, expectedRevision)
  1947. else if (mode === 'replace') await settings.replace(branded, section, expectedRevision)
  1948. else await settings.mutate(branded, section as SettingsPathOp[], expectedRevision)
  1949. } catch (error: unknown) {
  1950. return rejected(error)
  1951. }
  1952. const descriptor = settings.describe({ redactSecrets: true }).find(candidate => candidate.ns === branded)
  1953. if (descriptor === undefined) {
  1954. // The write committed but the namespace vanished before this read: only
  1955. // a concurrent registrant disposal can produce it.
  1956. return err(request, { code: 'internal', message: `settings namespace "${ns}" was disposed after the ${mode}`, details: {} })
  1957. }
  1958. return ok(request, namespaceView(descriptor))
  1959. }
  1960. return {
  1961. sessions: {
  1962. // Attached sessions summarize from memory; persisted-but-unattached (cold)
  1963. // sessions merge in from the persistence store so history survives restarts.
  1964. // Logs without a cwd are not served; every session records its project
  1965. // at create time.
  1966. async list(request) {
  1967. return ok(request, { items: await listVisibleSessionSummaries() })
  1968. },
  1969. async search(request, signal) {
  1970. const cancelled = () => err<{ items: SessionSearchItem[]; hasMore: boolean }>(request, {
  1971. code: 'cancelled',
  1972. message: 'session search was aborted',
  1973. details: {},
  1974. })
  1975. if (isAborted(signal)) return cancelled()
  1976. const sessionQuery = ctx.get('sessionQuery')
  1977. if (sessionQuery === undefined) {
  1978. return err(request, {
  1979. code: 'internal',
  1980. message: 'session search is unavailable: this deployment does not mount @deepseek-ai/dsh-session-query',
  1981. details: {},
  1982. })
  1983. }
  1984. try {
  1985. const visible = await listVisibleSessionSummaries(signal)
  1986. if (isAborted(signal)) return cancelled()
  1987. if (visible.length === 0) return ok(request, { items: [], hasMore: false })
  1988. const visibleIds = new Set(visible.map(item => item.sessionId))
  1989. const authorized: SessionSearchItem[] = []
  1990. const acceptedIds = new Set<SessionId>()
  1991. const seenCursors = new Set<SessionSearchCursor>()
  1992. let cursor: SessionSearchCursor | undefined
  1993. let providerCallCount = 0
  1994. let providerPageLimit = SESSION_SEARCH_RESULT_LIMIT
  1995. while (authorized.length <= SESSION_SEARCH_RESULT_LIMIT) {
  1996. if (isAborted(signal)) return cancelled()
  1997. if (providerCallCount >= SESSION_SEARCH_PROVIDER_CALL_LIMIT) {
  1998. throw new Error(
  1999. `session search provider exceeded the ${SESSION_SEARCH_PROVIDER_CALL_LIMIT}-call work budget`,
  2000. )
  2001. }
  2002. providerCallCount++
  2003. const requestedCursor = cursor
  2004. const requestedPageLimit = providerPageLimit
  2005. let page
  2006. try {
  2007. page = await sessionQuery.searchSessions({
  2008. query: request.payload.query,
  2009. eventFilters: [
  2010. { kind: 'type', values: ['user/message', 'assistant/message'] },
  2011. { kind: 'surface', values: ['current'] },
  2012. ],
  2013. limit: requestedPageLimit,
  2014. ...requestedCursor === undefined ? {} : { cursor: requestedCursor },
  2015. }, { signal })
  2016. } catch (error: unknown) {
  2017. if (isAborted(signal)) return cancelled()
  2018. if (
  2019. requestedCursor === undefined
  2020. && error instanceof SessionQueryError
  2021. && error.code === 'SESSION_QUERY_INVALID_LIMIT'
  2022. && requestedPageLimit > 1
  2023. ) {
  2024. providerPageLimit = Math.max(1, Math.floor(requestedPageLimit / 2))
  2025. continue
  2026. }
  2027. if (
  2028. requestedCursor !== undefined
  2029. && error instanceof SessionQueryError
  2030. && error.code === 'SESSION_QUERY_STALE_CURSOR'
  2031. ) {
  2032. authorized.length = 0
  2033. acceptedIds.clear()
  2034. seenCursors.clear()
  2035. cursor = undefined
  2036. continue
  2037. }
  2038. throw error
  2039. }
  2040. if (isAborted(signal)) return cancelled()
  2041. const providerItemCount = page.items.length
  2042. if (providerItemCount > requestedPageLimit) {
  2043. throw new Error(
  2044. `session search provider returned ${providerItemCount} items; maximum is ${requestedPageLimit}`,
  2045. )
  2046. }
  2047. // Host visibility is the authorization boundary. Consume the
  2048. // provider's globally ranked results rather than binding every
  2049. // visible id into one SQLite statement, then require each hit to
  2050. // name a visible session and a current message from that same
  2051. // session before emitting its snippet.
  2052. for (const hit of page.items) {
  2053. if (authorized.length > SESSION_SEARCH_RESULT_LIMIT) continue
  2054. if (
  2055. !visibleIds.has(hit.header.id)
  2056. || hit.bestMatch.sessionId !== hit.header.id
  2057. || hit.bestMatch.surface !== 'current'
  2058. || !MESSAGE_TYPES.has(hit.bestMatch.type)
  2059. || acceptedIds.has(hit.header.id)
  2060. ) continue
  2061. const snippet = truncateUnicodeCodePoints(
  2062. hit.bestMatch.snippet,
  2063. SESSION_SEARCH_SNIPPET_MAX_CODE_POINTS,
  2064. )
  2065. acceptedIds.add(hit.header.id)
  2066. authorized.push({
  2067. sessionId: hit.header.id,
  2068. snippet,
  2069. })
  2070. }
  2071. const nextCursor = page.nextCursor
  2072. if (nextCursor !== undefined) {
  2073. if (seenCursors.has(nextCursor)) {
  2074. throw new Error('session search provider repeated a continuation cursor')
  2075. }
  2076. seenCursors.add(nextCursor)
  2077. }
  2078. if (authorized.length > SESSION_SEARCH_RESULT_LIMIT || nextCursor === undefined) break
  2079. cursor = nextCursor
  2080. }
  2081. return ok(request, {
  2082. items: authorized.slice(0, SESSION_SEARCH_RESULT_LIMIT),
  2083. hasMore: authorized.length > SESSION_SEARCH_RESULT_LIMIT,
  2084. })
  2085. } catch (error: unknown) {
  2086. if (
  2087. isAborted(signal)
  2088. || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')
  2089. ) return cancelled()
  2090. // XXX: Redact provider details before exposing this gateway beyond
  2091. // its current single-user local deployment.
  2092. return err(request, {
  2093. code: 'internal',
  2094. message: `session search failed: ${String(error)}`,
  2095. details: {},
  2096. })
  2097. }
  2098. },
  2099. async create(request) {
  2100. const sessionId = request.payload.sessionId ?? `session-${randomUUID()}` as SessionId
  2101. let workspace: Workspace | undefined
  2102. if (request.payload.workspaceId !== undefined) {
  2103. workspace = ctx.workspaceRegistry.get(brandWorkspaceId(request.payload.workspaceId))
  2104. if (workspace === undefined) {
  2105. return err(request, {
  2106. code: 'workspace-not-found',
  2107. message: `workspace "${request.payload.workspaceId}" not found`,
  2108. details: { workspaceId: request.payload.workspaceId },
  2109. })
  2110. }
  2111. }
  2112. const cwd = workspace?.path ?? request.payload.cwd ?? defaults.cwd
  2113. const requestedPreset = request.payload.agentPreset
  2114. try {
  2115. await ensureSession(sessionId, cwd, request.payload.sessionId !== undefined, requestedPreset)
  2116. } catch (error: unknown) {
  2117. if (error instanceof AgentPresetConflict) {
  2118. return err(request, {
  2119. code: 'agent-preset-conflict',
  2120. message: error.message,
  2121. details: {
  2122. sessionId: error.sessionId,
  2123. requestedPreset: error.requestedPreset,
  2124. ...error.existingPreset === undefined ? {} : { existingPreset: error.existingPreset },
  2125. },
  2126. })
  2127. }
  2128. const refused = presetFailure(request, error)
  2129. if (refused !== undefined) return refused
  2130. if (error instanceof SessionCwdConflict) {
  2131. return err(request, {
  2132. code: 'session-conflict',
  2133. message: error.message,
  2134. details: {
  2135. sessionId: error.sessionId,
  2136. requestedCwd: error.requestedCwd,
  2137. ...error.existingCwd === undefined ? {} : { existingCwd: error.existingCwd },
  2138. },
  2139. })
  2140. }
  2141. if (error instanceof SubagentSessionOwnership) {
  2142. return err(request, subagentOwnershipError(error.sessionId))
  2143. }
  2144. return err(request, {
  2145. code: 'internal',
  2146. message: `failed to create session "${sessionId}": ${String(error)}`,
  2147. details: {},
  2148. })
  2149. }
  2150. if (workspace !== undefined) {
  2151. try {
  2152. await workspace.attachSession(sessionId)
  2153. } catch (error: unknown) {
  2154. return err(request, {
  2155. code: 'workspace-attach-failed',
  2156. message: `session "${sessionId}" was created but could not attach to workspace "${workspace.id}": ${String(error)}`,
  2157. details: { sessionId, workspaceId: workspace.id },
  2158. })
  2159. }
  2160. }
  2161. // Echo the composition the session RUNS so a client can label it
  2162. // without waiting for the next list refresh — the create is the commit
  2163. // point that knows it (a caller that named none gets the default).
  2164. // Resolved from the log for the same reason `sessionListFields()` is:
  2165. // this handler also adopts an already-live session, and one that
  2166. // switched while blank runs a preset its header no longer names, so
  2167. // echoing the header would contradict both the adoption this call just
  2168. // allowed and the row `session.list` serves for the same session.
  2169. const created = ctx.agents.get(sessionId)
  2170. const createdPreset = created === undefined ? undefined : resolveSessionPreset(created.session)
  2171. return ok(request, { sessionId, ...createdPreset === undefined ? {} : { agentPreset: createdPreset } })
  2172. },
  2173. async history(request) {
  2174. const { sessionId, beforeSeq, maxMessages } = request.payload
  2175. try {
  2176. const source = await historySourceFor(sessionId)
  2177. // Both awaits happen BEFORE the cut. Ensuring the recorded
  2178. // composition's standing mount is what registers its projection
  2179. // units, so a first cold read would otherwise serve a baseline
  2180. // missing every preset-owned key; and an attached session keeps
  2181. // appending, so awaiting between the two reads would pair events cut
  2182. // at N with a baseline folded to N+1.
  2183. const scope = await presenterScopeFor(sessionId, sourceSession(source))
  2184. const cut = historyCutOf(source, beforeSeq === undefined)
  2185. const page = historyPage(ctx, cut.events, beforeSeq, maxMessages, scope)
  2186. return ok(request, {
  2187. events: page.events,
  2188. hasMore: page.hasMore,
  2189. ...cut.projections === undefined ? {} : { projections: cut.projections },
  2190. })
  2191. } catch (error: unknown) {
  2192. if (error instanceof SessionNotFound) {
  2193. return err(request, { code: 'session-not-found', message: error.message, details: { sessionId } })
  2194. }
  2195. return err(request, {
  2196. code: 'internal',
  2197. message: `history unavailable for session "${sessionId}": ${String(error)}`,
  2198. details: {},
  2199. })
  2200. }
  2201. },
  2202. async models(request) {
  2203. const { sessionId } = request.payload
  2204. const found = await agentFor(sessionId)
  2205. if ('error' in found) return err(request, found.error)
  2206. const current = selectionFor(found.agent).current
  2207. const { groups, failures } = await buildModelCatalog(ctx)
  2208. const routable = routeServed(current.provider)
  2209. return ok(request, { current: { ...current }, routable, groups, failures })
  2210. },
  2211. async selectModel(request) {
  2212. const { sessionId, provider, model, reasoningEffort } = request.payload
  2213. const found = await agentFor(sessionId)
  2214. if ('error' in found) return err(request, found.error)
  2215. return serializeImageAdmission(found.agent, async () => {
  2216. try {
  2217. const resolved = await ctx.llm.resolveCallConfig({
  2218. provider,
  2219. model,
  2220. ...reasoningEffort === undefined
  2221. ? {}
  2222. : { reasoningEffort: ReasoningEffortId(reasoningEffort) },
  2223. })
  2224. const pendingImage = [...found.agent.inbox.nextTurn, ...found.agent.inbox.nextStep]
  2225. .some(message => contentHasImage(message.content))
  2226. if (pendingImage || messagesHaveImage(found.agent.session.deriveMessages())) {
  2227. const info = await ctx.llm.resolveModelInfo(resolved.provider, resolved.model)
  2228. if (info.inputModalities !== undefined && !info.inputModalities.includes('image')) {
  2229. return err(request, {
  2230. code: 'model-unavailable',
  2231. message: `Model "${resolved.model}" does not accept image input, but this session already contains images; select an image-capable model.`,
  2232. details: { provider, model },
  2233. })
  2234. }
  2235. }
  2236. const selected: ModelSelection = {
  2237. provider: resolved.provider,
  2238. model: resolved.model,
  2239. ...resolved.reasoningEffort === undefined
  2240. ? {}
  2241. : { reasoningEffort: resolved.reasoningEffort },
  2242. }
  2243. selectionFor(found.agent).current = selected
  2244. try {
  2245. await defaults.saveDefaultModelSelection?.(selected)
  2246. } catch (error: unknown) {
  2247. ctx.logger.warn(
  2248. `api-proxy: the model switch applies to this session but was not saved as the default: ${String(error)}`,
  2249. )
  2250. }
  2251. return ok(request, { selected: { ...selected } })
  2252. } catch (error: unknown) {
  2253. return err(request, {
  2254. code: 'model-unavailable',
  2255. message: error instanceof Error ? error.message : String(error),
  2256. details: { provider, model },
  2257. })
  2258. }
  2259. })
  2260. },
  2261. async rename(request) {
  2262. const { sessionId, title } = request.payload
  2263. const found = await agentFor(sessionId)
  2264. if ('error' in found) return err(request, found.error)
  2265. const titles = ctx.get('sessionTitle')
  2266. if (titles === undefined) {
  2267. return err(request, { code: 'internal', message: 'renaming is unavailable: this deployment mounts no session-title service', details: {} })
  2268. }
  2269. try {
  2270. const accepted = titles.rename(found.agent.session, title)
  2271. return ok(request, { title: accepted.title, seq: accepted.eventSeq })
  2272. } catch (error: unknown) {
  2273. // Only the input's fault maps to title-invalid (the message is
  2274. // product-user-visible in the rename dialog); liveness and disposal
  2275. // races are deployment trouble, not a bad title.
  2276. if (error instanceof SessionTitleInvalidError) {
  2277. return err(request, {
  2278. code: 'title-invalid',
  2279. message: error.message,
  2280. details: { sessionId },
  2281. })
  2282. }
  2283. return err(request, {
  2284. code: 'internal',
  2285. message: `failed to rename session "${sessionId}": ${String(error)}`,
  2286. details: {},
  2287. })
  2288. }
  2289. },
  2290. async fork(request) {
  2291. const { sessionId, atSeq } = request.payload
  2292. let source: SessionReadState
  2293. try {
  2294. source = await readSessionState(sessionId)
  2295. } catch (error: unknown) {
  2296. if (error instanceof SessionNotFound) {
  2297. return err(request, { code: 'session-not-found', message: error.message, details: { sessionId } })
  2298. }
  2299. return err(request, {
  2300. code: 'internal',
  2301. message: `fork source unavailable for session "${sessionId}": ${String(error)}`,
  2302. details: {},
  2303. })
  2304. }
  2305. const events = source.events
  2306. // An in-log anchor belongs to the turn containing it and must never
  2307. // clip backward to an earlier completed turn. Omitted and past-end
  2308. // anchors retain the last-completed-turn shortcut.
  2309. const lastSeq = events.at(-1)?.seq ?? -1
  2310. const anchoredBoundary = atSeq === undefined
  2311. ? undefined
  2312. : events.find(e => e.type === 'turn/end' && e.seq >= atSeq)
  2313. const boundary = anchoredBoundary
  2314. ?? (atSeq === undefined || atSeq > lastSeq
  2315. ? events.findLast(e => e.type === 'turn/end')
  2316. : undefined)
  2317. if (boundary === undefined) {
  2318. return err(request, {
  2319. code: 'fork-unavailable',
  2320. message: atSeq !== undefined && atSeq <= lastSeq
  2321. ? `session "${sessionId}" has not completed the turn containing event ${String(atSeq)}`
  2322. : `session "${sessionId}" has no completed turn to fork from`,
  2323. details: { sessionId },
  2324. })
  2325. }
  2326. // Extend the cut through trailing out-of-band appends (session/title,
  2327. // injections) up to the next turn/start: they are standalone events, so
  2328. // the seed stays balanced, and the child inherits a title generated
  2329. // right after the boundary turn.
  2330. let cut = boundary.seq + 1
  2331. while (cut < events.length && events[cut]?.type !== 'turn/start') cut++
  2332. let workspace: Workspace | undefined
  2333. try {
  2334. workspace = await forkWorkspace(source)
  2335. } catch (error: unknown) {
  2336. return err(request, {
  2337. code: 'internal',
  2338. message: `failed to resolve fork workspace for session "${sessionId}": ${String(error)}`,
  2339. details: {},
  2340. })
  2341. }
  2342. const childId = `session-${randomUUID()}` as SessionId
  2343. // The child inherits the parent's composition for the same reason a
  2344. // resumed session keeps its own: the seeded history was produced under
  2345. // those tools, and composing anything else would strand the tool calls
  2346. // it already carries. Now that no model-facing row sits in the host
  2347. // plane, composing nothing would leave the child with no tools at all.
  2348. const forkComposition = await composeAgent(resolveSessionPreset(source))
  2349. try {
  2350. await ctx.agents.create({
  2351. sessionId: childId,
  2352. seed: events.slice(0, cut),
  2353. meta: {
  2354. ...source.header.cwd === undefined ? {} : { cwd: source.header.cwd },
  2355. parentSession: source.id,
  2356. seedLength: cut,
  2357. ...forkComposition.agentPreset === undefined
  2358. ? {}
  2359. : { agentPreset: forkComposition.agentPreset },
  2360. },
  2361. agentOptions: agentOptions(),
  2362. setup: forkComposition.setup,
  2363. })
  2364. } catch (error: unknown) {
  2365. return err(request, {
  2366. code: 'internal',
  2367. message: `failed to fork session "${sessionId}": ${String(error)}`,
  2368. details: {},
  2369. })
  2370. }
  2371. // An ordinary source keeps its direct Workspace. A subagent source is
  2372. // not listed there, so its ordinary fork joins the nearest owning
  2373. // ancestor instead. The child is already published if attach fails.
  2374. if (workspace !== undefined) {
  2375. try {
  2376. await workspace.attachSession(childId)
  2377. } catch (error: unknown) {
  2378. return err(request, {
  2379. code: 'workspace-attach-failed',
  2380. message: `session "${childId}" was forked but could not attach to workspace "${workspace.id}": ${String(error)}`,
  2381. details: { sessionId: childId, workspaceId: workspace.id },
  2382. })
  2383. }
  2384. }
  2385. return ok(request, { sessionId: childId })
  2386. },
  2387. async prompt(request, signal) {
  2388. const { sessionId, mode, content, clientTimeZone } = request.payload
  2389. const canonicalTimeZone = clientTimeZone === undefined
  2390. ? undefined
  2391. : canonicalClientTimeZone(clientTimeZone)
  2392. if (clientTimeZone !== undefined && canonicalTimeZone === undefined) {
  2393. return err(request, {
  2394. code: 'invalid-time-zone',
  2395. message: 'clientTimeZone must be UTC or a valid IANA Area/Location name',
  2396. details: { value: clientTimeZone },
  2397. })
  2398. }
  2399. const resolved = await turnAgentFor<{ accepted: true }>(request, sessionId)
  2400. if ('refused' in resolved) return resolved.refused
  2401. const agent = resolved.agent
  2402. let parsed: ReturnType<typeof parseReferencedContent>
  2403. try {
  2404. parsed = parseReferencedContent(content)
  2405. } catch (error: unknown) {
  2406. return err(request, {
  2407. code: 'reference-invalid',
  2408. message: 'invalid session reference',
  2409. details: { reason: String(error) },
  2410. })
  2411. }
  2412. // Request identity and optional browser zone ride the exact durable user message.
  2413. const source: MessageSource = {
  2414. kind: 'user',
  2415. rpcId: request.rpcId,
  2416. ...(canonicalTimeZone === undefined ? {} : { clientTimeZone: canonicalTimeZone }),
  2417. }
  2418. const hasImage = parsed.content.some(part => part.type === 'image')
  2419. const admit = async (): Promise<RpcResponse<{ accepted: true }>> => {
  2420. try {
  2421. if (signal?.aborted === true) {
  2422. return err(request, {
  2423. code: 'cancelled',
  2424. message: 'prompt submission was aborted',
  2425. details: {},
  2426. })
  2427. }
  2428. if (hasImage) {
  2429. const current = selectionFor(agent).current
  2430. const modelInfo = await ctx.llm.resolveModelInfo(current.provider, current.model)
  2431. if (modelInfo.inputModalities !== undefined && !modelInfo.inputModalities.includes('image')) {
  2432. return err(request, {
  2433. code: 'attachment-error',
  2434. message: `Model "${current.model}" does not support image input.`,
  2435. details: { reason: 'MODEL_DOES_NOT_SUPPORT_IMAGES' },
  2436. })
  2437. }
  2438. }
  2439. let durable = await durablePromptContent(ctx, parsed.content)
  2440. let additionalContext: UserMessage | undefined
  2441. if (parsed.references.length > 0) {
  2442. const sessionReferences = ctx.get('sessionReferenceResolver')
  2443. if (sessionReferences === undefined) {
  2444. return err(request, {
  2445. code: 'reference-unavailable',
  2446. message: 'session reference capability unavailable',
  2447. details: { kind: 'session' },
  2448. })
  2449. }
  2450. try {
  2451. const prepared = await sessionReferences.prepare(agent, durable, parsed.references, signal)
  2452. durable = prepared.content
  2453. additionalContext = prepared.additionalContext
  2454. } catch (error: unknown) {
  2455. if (signal !== undefined && isAborted(signal)) {
  2456. return err(request, {
  2457. code: 'cancelled',
  2458. message: 'session reference preparation was aborted',
  2459. details: {},
  2460. })
  2461. }
  2462. return err(request, {
  2463. code: 'reference-failed',
  2464. message: 'session reference preparation failed',
  2465. details: { reason: String(error) },
  2466. })
  2467. }
  2468. }
  2469. if (signal !== undefined && isAborted(signal)) {
  2470. return err(request, {
  2471. code: 'cancelled',
  2472. message: 'prompt submission was aborted',
  2473. details: {},
  2474. })
  2475. }
  2476. const message: UserMessage = createUserMessage({ content: durable, source })
  2477. deliverPrompt(ctx, agent, mode, message, additionalContext, preparedPromptOwnership)
  2478. } catch (error: unknown) {
  2479. if (error instanceof AttachmentError) {
  2480. return err(request, {
  2481. code: 'attachment-error',
  2482. message: error.message,
  2483. details: { reason: error.code },
  2484. })
  2485. }
  2486. return err(request, {
  2487. code: 'agent-busy',
  2488. message: 'prompt rejected',
  2489. details: { reason: String(error) },
  2490. })
  2491. }
  2492. return ok(request, { accepted: true as const })
  2493. }
  2494. return hasImage ? serializeImageAdmission(agent, admit) : admit()
  2495. },
  2496. async attachment(request) {
  2497. const { sessionId, attachmentId } = request.payload
  2498. let state: SessionReadState
  2499. try {
  2500. state = await readSessionState(sessionId)
  2501. } catch (error: unknown) {
  2502. if (error instanceof SessionNotFound) {
  2503. return err(request, {
  2504. code: 'session-not-found',
  2505. message: error.message,
  2506. details: { sessionId },
  2507. })
  2508. }
  2509. return err(request, {
  2510. code: 'internal',
  2511. message: `attachment authorization unavailable for session "${sessionId}": ${String(error)}`,
  2512. details: {},
  2513. })
  2514. }
  2515. const ref = referencedImage(state.events, String(attachmentId))
  2516. if (ref === undefined) {
  2517. return err(request, {
  2518. code: 'attachment-error',
  2519. message: 'Image is not referenced by this session.',
  2520. details: { reason: 'ATTACHMENT_NOT_REFERENCED' },
  2521. })
  2522. }
  2523. try {
  2524. const stored = await ctx.attachments.readImage(ref)
  2525. return ok(request, {
  2526. attachment: stored.ref,
  2527. data: Buffer.from(stored.data).toString('base64'),
  2528. })
  2529. } catch (error: unknown) {
  2530. if (error instanceof AttachmentError) {
  2531. return err(request, {
  2532. code: 'attachment-error',
  2533. message: error.message,
  2534. details: { reason: error.code },
  2535. })
  2536. }
  2537. return err(request, {
  2538. code: 'internal',
  2539. message: 'Unable to read image attachment.',
  2540. details: {},
  2541. })
  2542. }
  2543. },
  2544. updateQueue(request) {
  2545. const { sessionId, itemId, action } = request.payload
  2546. if (action.kind === 'edit' && action.content.some(block => block.type !== 'text')) {
  2547. return Promise.resolve(err(request, {
  2548. code: 'attachment-error',
  2549. message: 'queue edits accept text content only',
  2550. details: { reason: 'QUEUE_EDIT_NON_TEXT' },
  2551. }))
  2552. }
  2553. const agent = ctx.agents.get(sessionId)
  2554. if (agent !== undefined && hasSubagentOwner(agent.session, agent)) {
  2555. return Promise.resolve(err(request, subagentOwnershipError(sessionId)))
  2556. }
  2557. if (agent === undefined) {
  2558. return Promise.resolve(err(request, {
  2559. code: 'queue-item-not-found',
  2560. message: 'queued item is no longer pending',
  2561. details: { itemId },
  2562. }))
  2563. }
  2564. const target = agent.inbox.nextTurn.some(message => message.id === itemId)
  2565. ? 'next-turn'
  2566. : agent.inbox.nextStep.some(message => message.id === itemId) ? 'next-step' : undefined
  2567. const message = target === undefined
  2568. ? undefined
  2569. : (target === 'next-turn' ? agent.inbox.nextTurn : agent.inbox.nextStep)
  2570. .find(candidate => candidate.id === itemId)
  2571. if (target === undefined || message === undefined) {
  2572. return Promise.resolve(err(request, {
  2573. code: 'queue-item-not-found',
  2574. message: 'queued item is no longer pending',
  2575. details: { itemId },
  2576. }))
  2577. }
  2578. if (action.kind === 'steer' && (target !== 'next-turn' || agent.status !== 'running')) {
  2579. return Promise.resolve(err(request, {
  2580. code: 'steer-unavailable',
  2581. message: 'current turn no longer accepts steering',
  2582. details: { itemId },
  2583. }))
  2584. }
  2585. if (action.kind === 'edit') {
  2586. agent.inbox.replace(itemId, freezeMessage({ ...message, content: action.content }))
  2587. } else {
  2588. if (action.kind === 'steer') preparedPromptOwnership.relocating.add(itemId)
  2589. try {
  2590. agent.inbox.remove(itemId)
  2591. if (action.kind === 'steer') agent.steer(message)
  2592. } catch (error: unknown) {
  2593. preparedPromptOwnership.cleanups.get(itemId)?.()
  2594. throw error
  2595. } finally {
  2596. preparedPromptOwnership.relocating.delete(itemId)
  2597. }
  2598. }
  2599. return Promise.resolve(ok(request, { accepted: true as const }))
  2600. },
  2601. cancel(request) {
  2602. const { sessionId } = request.payload
  2603. const agent = ctx.agents.get(sessionId)
  2604. if (agent === undefined) {
  2605. return Promise.resolve(err(request, {
  2606. code: 'session-not-found',
  2607. message: `session "${sessionId}" not found (not attached)`,
  2608. details: { sessionId },
  2609. }))
  2610. }
  2611. if (hasSubagentOwner(agent.session, agent)) {
  2612. return Promise.resolve(err(request, subagentOwnershipError(sessionId)))
  2613. }
  2614. agent.cancel({ kind: 'user' }, { keepInbox: true })
  2615. return Promise.resolve(ok(request, { accepted: true as const }))
  2616. },
  2617. },
  2618. subagents: {
  2619. async list(request, signal) {
  2620. try {
  2621. const entries = await ctx.subagents.listChildren(request.payload.parentSessionId, signal)
  2622. return ok(request, {
  2623. entries: entries.map(entry => entry.kind === 'child'
  2624. ? {
  2625. ...entry,
  2626. activity: ctx.agents.get(entry.id)?.status === 'running' ? 'running' : 'inactive',
  2627. }
  2628. : entry),
  2629. parentAvailable: ctx.agents.get(request.payload.parentSessionId) !== undefined,
  2630. })
  2631. } catch (error: unknown) {
  2632. if (signal?.aborted || (error instanceof SubagentError && error.code === 'CANCELLED')) {
  2633. return err(request, {
  2634. code: 'cancelled',
  2635. message: 'subagent catalog read was cancelled',
  2636. details: {},
  2637. })
  2638. }
  2639. if (error instanceof SubagentError && error.code === 'SUBAGENT_CONTROL_PROJECTIONS_UNAVAILABLE') {
  2640. return err(request, projectionsUnavailableError())
  2641. }
  2642. return err(request, {
  2643. code: 'internal',
  2644. message: 'subagent catalog read failed',
  2645. details: {},
  2646. })
  2647. }
  2648. },
  2649. async history(request, signal) {
  2650. const {
  2651. parentSessionId, childSessionId, mode, beforeSeq, maxMessages,
  2652. } = request.payload
  2653. const verified = await catalogChild(ctx, {
  2654. parentSessionId, childSessionId, mode,
  2655. }, signal)
  2656. if (verified.error !== undefined) return err(request, verified.error)
  2657. // The generic-history data plane: an attached child serves its
  2658. // in-memory snapshot and the registry's live watermark projections; a
  2659. // cold child is one persistence inspection plus a detached fold.
  2660. let header: SessionHeader
  2661. let events: SessionEvent[]
  2662. let projections: SessionProjectionsBlock | undefined
  2663. const attached = ctx.sessions.get(childSessionId)
  2664. if (attached !== undefined) {
  2665. header = attached.header
  2666. events = [...attached.events]
  2667. projections = beforeSeq === undefined
  2668. ? subagentHistoryProjections(ctx, childSessionId, () => projectionsFor(ctx, attached))
  2669. : undefined
  2670. } else {
  2671. try {
  2672. const inspected = await inspectServable(childSessionId)
  2673. header = inspected.meta
  2674. events = inspected.events
  2675. projections = beforeSeq === undefined
  2676. ? subagentHistoryProjections(ctx, childSessionId, () => detachedProjectionsFor(ctx, inspected.events))
  2677. : undefined
  2678. } catch (error: unknown) {
  2679. if (signal?.aborted) {
  2680. return err(request, {
  2681. code: 'cancelled',
  2682. message: 'subagent history read was cancelled',
  2683. details: {},
  2684. })
  2685. }
  2686. if (error instanceof SessionNotFound) {
  2687. return err(request, {
  2688. code: 'subagent-not-found',
  2689. message: 'subagent disappeared during history read',
  2690. details: { parentSessionId, childSessionId },
  2691. })
  2692. }
  2693. return err(request, {
  2694. code: 'internal',
  2695. message: 'subagent history read failed',
  2696. details: {},
  2697. })
  2698. }
  2699. }
  2700. if (signal?.aborted) {
  2701. return err(request, {
  2702. code: 'cancelled',
  2703. message: 'subagent history read was cancelled',
  2704. details: {},
  2705. })
  2706. }
  2707. if (header.parentSession !== parentSessionId) {
  2708. return err(request, {
  2709. code: 'subagent-unauthorized',
  2710. message: 'subagent parent changed during history read',
  2711. details: { childSessionId },
  2712. })
  2713. }
  2714. const page = historyPage(ctx, events, beforeSeq, maxMessages)
  2715. return ok(request, { ...page, ...projections === undefined ? {} : { projections } })
  2716. },
  2717. async prompt(request, signal) {
  2718. const { parentSessionId, childSessionId, content, clientTimeZone } = request.payload
  2719. const canonicalTimeZone = clientTimeZone === undefined
  2720. ? undefined
  2721. : canonicalClientTimeZone(clientTimeZone)
  2722. if (clientTimeZone !== undefined && canonicalTimeZone === undefined) {
  2723. return err(request, {
  2724. code: 'invalid-time-zone',
  2725. message: 'clientTimeZone must be UTC or a valid IANA Area/Location name',
  2726. details: { value: clientTimeZone },
  2727. })
  2728. }
  2729. const parent = ctx.agents.get(parentSessionId)
  2730. if (parent === undefined) {
  2731. return err(request, {
  2732. code: 'subagent-parent-unavailable',
  2733. message: `parent session "${parentSessionId}" is not live`,
  2734. details: { parentSessionId },
  2735. })
  2736. }
  2737. const verified = await catalogChild(ctx, {
  2738. parentSessionId, childSessionId, mode: 'continuable',
  2739. }, signal)
  2740. if (verified.error !== undefined) return err(request, verified.error)
  2741. try {
  2742. const messageId = await ctx.subagents.followup(parent, childSessionId, content, {
  2743. source: {
  2744. kind: 'user',
  2745. rpcId: request.rpcId,
  2746. ...(canonicalTimeZone === undefined ? {} : { clientTimeZone: canonicalTimeZone }),
  2747. },
  2748. signal,
  2749. })
  2750. return ok(request, { messageId })
  2751. } catch (error: unknown) {
  2752. return subagentPromptError(request, error, signal)
  2753. }
  2754. },
  2755. // Deliberately no catalog, history, persistence, or parent Agent lookup:
  2756. // the core primitive alone authorizes the durable address against the
  2757. // live Activation, which is what keeps a live child interruptible while
  2758. // its parent Agent is offline. Absent targets are accepted no-ops there.
  2759. interrupt(request) {
  2760. const { parentSessionId, childSessionId } = request.payload
  2761. try {
  2762. ctx.subagents.interrupt(childSessionId, { kind: 'user', parentSessionId })
  2763. } catch (error: unknown) {
  2764. if (error instanceof SubagentError && error.code === 'UNAUTHORIZED') {
  2765. return Promise.resolve(err(request, {
  2766. code: 'subagent-unauthorized',
  2767. message: 'subagent does not belong to this parent',
  2768. details: { childSessionId },
  2769. }))
  2770. }
  2771. return Promise.resolve(err(request, {
  2772. code: 'internal',
  2773. message: 'subagent interrupt failed',
  2774. details: {},
  2775. }))
  2776. }
  2777. return Promise.resolve(ok(request, { accepted: true as const }))
  2778. },
  2779. },
  2780. workspace: {
  2781. list(request) {
  2782. return Promise.resolve(ok(request, {
  2783. items: ctx.workspaceRegistry.list().map(workspaceView),
  2784. archivedSessionIds: [...ctx.workspaceRegistry.archivedSessionIds],
  2785. }))
  2786. },
  2787. async create(request) {
  2788. const { path } = request.payload
  2789. try {
  2790. const { workspace, created } = await ensureWorkspace(path)
  2791. return ok(request, { workspace: workspaceView(workspace), created })
  2792. } catch (error: unknown) {
  2793. // The registry rejects a path that does not resolve to an existing
  2794. // directory (realpath ENOENT / not-a-directory) — the business
  2795. // error of the typed-path flow, surfaced as a validation failure.
  2796. return err(request, {
  2797. code: 'workspace-invalid-path',
  2798. message: `cannot create a workspace at "${path}": ${error instanceof Error ? error.message : String(error)}`,
  2799. details: { path },
  2800. })
  2801. }
  2802. },
  2803. async rename(request) {
  2804. const { payload } = request
  2805. const workspace = ctx.workspaceRegistry.get(brandWorkspaceId(payload.workspaceId))
  2806. if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
  2807. const title = payload.title.trim()
  2808. // Uniqueness AND the same-title no-op both ride the create chain so
  2809. // they observe the state left by earlier queued renames — checked
  2810. // up front, a queued A→A could report success while an earlier A→B
  2811. // still lands afterwards.
  2812. const operation = workspaceCreationChain.then(async () => {
  2813. if (title === workspace.title) return
  2814. if (ctx.workspaceRegistry.list().some(other => other.id !== workspace.id && other.title === title)) {
  2815. throw new WorkspaceNameConflictError(title)
  2816. }
  2817. await workspace.setTitle(title)
  2818. })
  2819. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  2820. try {
  2821. await operation
  2822. } catch (error: unknown) {
  2823. if (error instanceof WorkspaceNameConflictError) {
  2824. return err(request, {
  2825. code: 'workspace-name-conflict',
  2826. message: error.message,
  2827. details: { name: error.workspaceName },
  2828. })
  2829. }
  2830. throw error
  2831. }
  2832. return ok(request, { workspace: workspaceView(workspace) })
  2833. },
  2834. async delete(request) {
  2835. const { workspaceId } = request.payload
  2836. const operation = workspaceCreationChain.then(() =>
  2837. ctx.workspaceRegistry.delete(brandWorkspaceId(workspaceId)))
  2838. workspaceCreationChain = operation.then(() => undefined, () => undefined)
  2839. if (!await operation) return workspaceNotFound(request, workspaceId)
  2840. return ok(request, { deleted: true as const })
  2841. },
  2842. async insertBefore(request) {
  2843. const { workspaceId, beforeWorkspaceId } = request.payload
  2844. try {
  2845. const workspaceIds = await ctx.workspaceRegistry.insertBefore(
  2846. brandWorkspaceId(workspaceId),
  2847. beforeWorkspaceId === undefined ? undefined : brandWorkspaceId(beforeWorkspaceId),
  2848. )
  2849. return ok(request, { workspaceIds: [...workspaceIds] })
  2850. } catch (error: unknown) {
  2851. if (!(error instanceof WorkspaceOrderInvalidError)) throw error
  2852. return workspaceNotFound(request, error.workspaceId)
  2853. }
  2854. },
  2855. async insertSessionBefore(request) {
  2856. const { payload } = request
  2857. const workspace = ctx.workspaceRegistry.get(brandWorkspaceId(payload.workspaceId))
  2858. if (workspace === undefined) return workspaceNotFound(request, payload.workspaceId)
  2859. try {
  2860. await workspace.insertSessionBefore(payload.sessionId, payload.beforeSessionId)
  2861. } catch (error: unknown) {
  2862. // Only the entity's unaccounted-id rejection is the business code;
  2863. // storage/durability failures propagate as internal errors.
  2864. if (!(error instanceof WorkspaceMoveInvalidError)) throw error
  2865. return err(request, {
  2866. code: 'workspace-move-invalid',
  2867. message: error.message,
  2868. details: {
  2869. workspaceId: payload.workspaceId,
  2870. sessionId: payload.sessionId,
  2871. ...payload.beforeSessionId === undefined ? {} : { beforeSessionId: payload.beforeSessionId },
  2872. },
  2873. })
  2874. }
  2875. return ok(request, { workspace: workspaceView(workspace) })
  2876. },
  2877. async archiveSession(request) {
  2878. const { sessionId } = request.payload
  2879. try {
  2880. await ctx.workspaceRegistry.archiveSession(sessionId)
  2881. } catch (error: unknown) {
  2882. // Only the registry's unknown-session rejection is the business
  2883. // code; storage/durability failures propagate as internal errors.
  2884. if (!(error instanceof WorkspaceUnknownSessionError)) throw error
  2885. return err(request, {
  2886. code: 'session-not-found',
  2887. message: error.message,
  2888. details: { sessionId },
  2889. })
  2890. }
  2891. return ok(request, { archivedSessionIds: [...ctx.workspaceRegistry.archivedSessionIds] })
  2892. },
  2893. },
  2894. host: {
  2895. describe(request) {
  2896. // TODO: version should read apps/cli's package.json; placeholder for now.
  2897. const selection = defaults.defaultModelSelection()
  2898. return Promise.resolve(ok(request, {
  2899. version: '0.0.1',
  2900. // Same source as session.create's fallback: the UI's default project
  2901. // must match where an unspecified-cwd session actually lands.
  2902. cwd: defaults.cwd,
  2903. // Read live for the same reason: this is what the NEXT session will
  2904. // start from, so a saved default has to be what it reports.
  2905. provider: selection.provider,
  2906. model: selection.model,
  2907. attachedSessions: ctx.agents.list().length,
  2908. canOpenPath: canOpenPaths(),
  2909. }))
  2910. },
  2911. async pickDirectory(request, signal) {
  2912. const capability = ctx.directoryPicker.capability()
  2913. if (capability.kind !== 'native') {
  2914. return err(request, {
  2915. code: 'directory-picker-unavailable',
  2916. message: `host.pickDirectory needs the native capability; the composed picker serves "${capability.kind}"`,
  2917. details: { capability: capability.kind },
  2918. })
  2919. }
  2920. try {
  2921. const path = await capability.pick(signal)
  2922. return ok(request, { path })
  2923. } catch (error: unknown) {
  2924. if (signal.aborted) {
  2925. return err(request, {
  2926. code: 'cancelled',
  2927. message: 'directory picker was aborted',
  2928. details: {},
  2929. })
  2930. }
  2931. return err(request, {
  2932. code: 'internal',
  2933. message: `directory picker failed: ${error instanceof Error ? error.message : String(error)}`,
  2934. details: {},
  2935. })
  2936. }
  2937. },
  2938. async listDirectory(request, signal) {
  2939. const capability = ctx.directoryPicker.capability()
  2940. if (capability.kind !== 'browse') {
  2941. return err(request, {
  2942. code: 'directory-picker-unavailable',
  2943. message: `host.listDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
  2944. details: { capability: capability.kind },
  2945. })
  2946. }
  2947. try {
  2948. // The carrier's signal follows the caller: a disconnect or timeout
  2949. // stops the backend's directory scan instead of outliving it.
  2950. return ok(request, await capability.list(request.payload.path, signal))
  2951. } catch (error: unknown) {
  2952. // An abort is the caller's own timeout/disconnect, not a server
  2953. // failure — same code pickDirectory and command.execute report.
  2954. if (signal.aborted) {
  2955. return err(request, { code: 'cancelled', message: 'directory listing was aborted', details: {} })
  2956. }
  2957. return err(request, directoryError(error))
  2958. }
  2959. },
  2960. async createDirectory(request) {
  2961. const capability = ctx.directoryPicker.capability()
  2962. if (capability.kind !== 'browse') {
  2963. return err(request, {
  2964. code: 'directory-picker-unavailable',
  2965. message: `host.createDirectory needs the browse capability; the composed picker serves "${capability.kind}"`,
  2966. details: { capability: capability.kind },
  2967. })
  2968. }
  2969. try {
  2970. return ok(request, { path: await capability.createDirectory(request.payload.path, request.payload.name) })
  2971. } catch (error: unknown) {
  2972. return err(request, directoryError(error))
  2973. }
  2974. },
  2975. async openPath(request, signal) {
  2976. return openPath(request, request.payload.path, signal)
  2977. },
  2978. },
  2979. goals: {
  2980. // Mutations only — the read side is the 'goal' session projection.
  2981. // Every verb resolves the session's agent (agentFor: implicit cold
  2982. // resume, the command.* precedent) and acknowledges with the new CAS
  2983. // ref; the committed goal/change event carries the whole value to every
  2984. // client through the projection frames.
  2985. async create(request) {
  2986. const { objective, maxGoalRounds } = request.payload
  2987. return mutateGoal(request, (goals, agent) => goals.create(agent, {
  2988. objective,
  2989. ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
  2990. }))
  2991. },
  2992. async edit(request) {
  2993. const { ref, objective, maxGoalRounds } = request.payload
  2994. return mutateGoal(request, (goals, agent) => goals.edit(agent, ref, {
  2995. ...(objective !== undefined ? { objective } : {}),
  2996. ...(maxGoalRounds !== undefined ? { maxGoalRounds } : {}),
  2997. }))
  2998. },
  2999. async pause(request) {
  3000. return mutateGoal(request, (goals, agent) => goals.pause(agent, request.payload.ref))
  3001. },
  3002. async resume(request) {
  3003. return mutateGoal(request, (goals, agent) => goals.resume(agent, request.payload.ref))
  3004. },
  3005. async complete(request) {
  3006. return mutateGoal(request, (goals, agent) => goals.complete(agent, request.payload.ref))
  3007. },
  3008. async clear(request) {
  3009. const found = await agentFor(request.payload.sessionId)
  3010. if ('error' in found) return err(request, found.error)
  3011. const goals = goalServiceFor(found.agent)
  3012. if ('error' in goals) return err(request, goals.error)
  3013. try {
  3014. goals.clear(found.agent, request.payload.ref)
  3015. return ok(request, { cleared: true as const })
  3016. } catch (error: unknown) {
  3017. return goalError(request, error)
  3018. }
  3019. },
  3020. },
  3021. agentPresets: {
  3022. // A deployment with no roster answers with an empty list rather than an
  3023. // error: composing no presets is a valid deployment, and the browser
  3024. // simply offers no choice.
  3025. async list(request) {
  3026. const presets = ctx.get('agentPresets')
  3027. if (presets === undefined) return ok(request, { presets: [], authorable: false, hasDocument: false })
  3028. const defaultId = presets.defaultId
  3029. return ok(request, {
  3030. presets: (await presets.list()).map(preset => ({
  3031. id: preset.id,
  3032. trust: preset.trust,
  3033. isDefault: preset.id === defaultId,
  3034. ...preset.name === undefined ? {} : { name: preset.name },
  3035. ...preset.description === undefined ? {} : { description: preset.description },
  3036. ...preset.broken === undefined ? {} : { broken: preset.broken },
  3037. })),
  3038. authorable: presets.authorable,
  3039. hasDocument: canOpenPaths(),
  3040. })
  3041. },
  3042. // Recomposing is limited to a blank session because a started
  3043. // conversation's history was produced under its preset's tools; the
  3044. // agent and the session survive, only the composition is swapped.
  3045. async select(request) {
  3046. const { sessionId, agentPreset } = request.payload
  3047. const presets = ctx.get('agentPresets')
  3048. if (presets === undefined) {
  3049. return err(request, {
  3050. code: 'agent-preset-not-found',
  3051. message: 'this deployment composes no agent presets',
  3052. details: { agentPreset, available: [] },
  3053. })
  3054. }
  3055. const found = await agentFor(sessionId)
  3056. if ('error' in found) return err(request, found.error)
  3057. const { agent } = found
  3058. const swap = async (): Promise<RpcResponse<{ agentPreset: string }>> => {
  3059. // Re-read inside the queue: an earlier switch may have run, and a
  3060. // conversation may have started, since this request arrived.
  3061. if (!sessionBlank(agent.session)) {
  3062. return err(request, {
  3063. code: 'agent-preset-locked',
  3064. message: `session "${sessionId}" has already started; its agent preset is fixed`,
  3065. details: { sessionId, agentPreset },
  3066. })
  3067. }
  3068. try {
  3069. const preset = await presets.recompose(agent.ctx, agentPreset)
  3070. // Recorded only after the swap committed: the log states what the
  3071. // agent runs, and a rejected mount leaves the previous composition.
  3072. agent.session.append('agent-preset/selected', { agentPreset: preset.id })
  3073. return ok(request, { agentPreset: preset.id })
  3074. } catch (error: unknown) {
  3075. const refused = presetFailure(request, error)
  3076. if (refused !== undefined) return refused
  3077. return err(request, {
  3078. code: 'internal',
  3079. message: `failed to select agent preset "${agentPreset}": ${String(error)}`,
  3080. details: {},
  3081. })
  3082. }
  3083. }
  3084. const queued = presetSwitches.get(sessionId) ?? Promise.resolve()
  3085. const turn = queued.then(swap)
  3086. presetSwitches.set(sessionId, turn.catch(() => undefined))
  3087. try {
  3088. return await turn
  3089. } finally {
  3090. if (presetSwitches.get(sessionId) === turn) presetSwitches.delete(sessionId)
  3091. }
  3092. },
  3093. // Authoring is privileged (see PRIVILEGED_METHODS in dsh-client-connection):
  3094. // a composition names the plugins a session runs, so reading one is
  3095. // reconnaissance, and copy/remove/openDocument manage the roster and
  3096. // drive the host desktop.
  3097. async read(request) {
  3098. const { agentPreset } = request.payload
  3099. const presets = ctx.get('agentPresets')
  3100. if (presets === undefined) return err(request, noRoster(agentPreset))
  3101. try {
  3102. const preset = await presets.resolve(agentPreset)
  3103. return ok(request, {
  3104. agentPreset: preset.id,
  3105. trust: preset.trust,
  3106. content: await presets.read(preset.id),
  3107. ...preset.name === undefined ? {} : { name: preset.name },
  3108. ...preset.description === undefined ? {} : { description: preset.description },
  3109. })
  3110. } catch (error: unknown) {
  3111. return err(request, presetError(agentPreset, error))
  3112. }
  3113. },
  3114. async copy(request) {
  3115. const { from, agentPreset, name } = request.payload
  3116. const presets = ctx.get('agentPresets')
  3117. if (presets === undefined) return err(request, noRoster(agentPreset))
  3118. try {
  3119. await presets.copy(from, agentPreset, name)
  3120. return ok(request, { agentPreset })
  3121. } catch (error: unknown) {
  3122. return err(request, presetError(agentPreset, error))
  3123. }
  3124. },
  3125. async openDocument(request, signal) {
  3126. const { agentPreset } = request.payload
  3127. const presets = ctx.get('agentPresets')
  3128. if (presets === undefined) return err(request, noRoster(agentPreset))
  3129. try {
  3130. const preset = await presets.resolve(agentPreset)
  3131. // Same line as copy/remove draw: the shipped install is not the
  3132. // user's to manage, and pointing an editor into it invites edits an
  3133. // upgrade will silently overwrite.
  3134. if (preset.trust !== 'user') {
  3135. throw new PresetNotWritableError(preset.id, 'it ships with the deployment')
  3136. }
  3137. // The id resolved against the Host's own roots is what selects the
  3138. // directory — no browser payload carries a path in either direction
  3139. // unless the deployment has no opener to hand it to.
  3140. const directory = dirname(preset.path)
  3141. if (!canOpenPaths()) return ok(request, { opened: false as const, path: directory })
  3142. return await openPath(request, directory, signal)
  3143. } catch (error: unknown) {
  3144. return err(request, presetError(agentPreset, error))
  3145. }
  3146. },
  3147. async remove(request) {
  3148. const { agentPreset } = request.payload
  3149. const presets = ctx.get('agentPresets')
  3150. if (presets === undefined) return err(request, noRoster(agentPreset))
  3151. try {
  3152. await presets.remove(agentPreset)
  3153. return ok(request, {})
  3154. } catch (error: unknown) {
  3155. return err(request, presetError(agentPreset, error))
  3156. }
  3157. },
  3158. },
  3159. skills: {
  3160. // Skill lookup never creates or resumes an agent: the session address
  3161. // resolves to a canonical cwd from the host-resident session header, and
  3162. // the view scope is the live agent or the preset's standing key.
  3163. async list(request) {
  3164. const { sessionId } = request.payload
  3165. const session = ctx.sessions.get(sessionId)
  3166. if (session === undefined) {
  3167. return err(request, {
  3168. code: 'session-not-found',
  3169. message: `session "${sessionId}" not found (not attached)`,
  3170. details: { sessionId },
  3171. })
  3172. }
  3173. if (session.header.cwd === undefined) {
  3174. // Every served session records its project at create time; a
  3175. // cwd-less header is a pre-project legacy log (not served).
  3176. return err(request, { code: 'internal', message: `session "${sessionId}" has no project cwd`, details: {} })
  3177. }
  3178. const cwd = session.header.cwd
  3179. // The host registry is layered per scope and serves every session. A
  3180. // composition may still realm-mount its own registry instead; that
  3181. // instance is invisible to host contexts, so address it through the
  3182. // live agent (`agents.get` keeps the no-side-effect stance above).
  3183. const live = ctx.agents.get(sessionId)
  3184. const presets = ctx.get('agentPresets')
  3185. const scoped = live === undefined ? undefined : presets?.serviceFor(live, 'skills')
  3186. // Same stance as the commands domain: a missing service means no
  3187. // composition mounts dsh-skill, not an empty catalog. `ctx.get` also
  3188. // keeps this handler independent of the gateway plugin's inject list
  3189. // (an undeclared `ctx.skills` property read fails the reflect proxy).
  3190. const skillRegistry = scoped ?? ctx.get('skills')
  3191. if (skillRegistry === undefined) {
  3192. return err(request, { code: 'internal', message: 'skill registry is absent: neither this session\'s agent preset nor the host composition mounts @deepseek-ai/dsh-skill', details: {} })
  3193. }
  3194. // The scope presenters resolve in — the live agent, else the recorded
  3195. // preset's standing key, else the global layer — so a cold session's
  3196. // '/' popup lists the catalog its composition actually serves.
  3197. const scope = await presenterScopeFor(sessionId, session)
  3198. try {
  3199. const skills = (await skillRegistry.list({ cwd, scope })).filter(isUserInvocable)
  3200. return ok(request, {
  3201. skills: skills.map(skill => ({
  3202. name: skill.name,
  3203. description: skill.description,
  3204. ...skill.whenToUse === undefined ? {} : { whenToUse: skill.whenToUse },
  3205. modelInvocable: skill.invocation.modelInvocable,
  3206. })),
  3207. })
  3208. } catch (error: unknown) {
  3209. return err(request, { code: 'internal', message: `skill listing failed: ${String(error)}`, details: {} })
  3210. }
  3211. },
  3212. },
  3213. settings: {
  3214. describe(request) {
  3215. const settings = ctx.get('settings')
  3216. if (settings === undefined) return Promise.resolve(err(request, settingsAbsent()))
  3217. return Promise.resolve(ok(request, {
  3218. writable: settings.writable,
  3219. hasDocument: settings.documentPath !== undefined,
  3220. namespaces: settings.describe({ redactSecrets: true }).map(namespaceView),
  3221. }))
  3222. },
  3223. async openDocument(request, signal) {
  3224. const settings = ctx.get('settings')
  3225. if (settings === undefined) return err(request, settingsAbsent())
  3226. if (isAborted(signal)) {
  3227. return err(request, {
  3228. code: 'cancelled',
  3229. message: 'settings document open was aborted',
  3230. details: {},
  3231. })
  3232. }
  3233. let path: string | undefined
  3234. try {
  3235. path = await settings.prepareDocument()
  3236. } catch (error: unknown) {
  3237. if (isAborted(signal)) {
  3238. return err(request, {
  3239. code: 'cancelled',
  3240. message: 'settings document preparation was aborted',
  3241. details: {},
  3242. })
  3243. }
  3244. return err(request, {
  3245. code: 'internal',
  3246. message: `settings document preparation failed: ${error instanceof Error ? error.message : String(error)}`,
  3247. details: {},
  3248. })
  3249. }
  3250. if (path === undefined) {
  3251. return err(request, {
  3252. code: 'internal',
  3253. message: 'settings provider has no local document to open',
  3254. details: {},
  3255. })
  3256. }
  3257. if (isAborted(signal)) {
  3258. return err(request, {
  3259. code: 'cancelled',
  3260. message: 'settings document open was aborted',
  3261. details: {},
  3262. })
  3263. }
  3264. return openTextFile(request, path, signal)
  3265. },
  3266. update: request => settingsWrite(request, request.payload.ns, 'update', request.payload.patch, request.payload.expectedRevision),
  3267. replace: request => settingsWrite(request, request.payload.ns, 'replace', request.payload.section, request.payload.expectedRevision),
  3268. mutate: request => settingsWrite(request, request.payload.ns, 'mutate', request.payload.ops, request.payload.expectedRevision),
  3269. },
  3270. credentials: {
  3271. async describe(request) {
  3272. const credentials = ctx.get('credentials')
  3273. if (credentials === undefined) return err(request, credentialsAbsent())
  3274. const entries = await Promise.all(request.payload.refs.map(async (ref) => {
  3275. const info = await credentials.describe(credentialRef(ref))
  3276. const view: CredentialView = {
  3277. configured: info.configured,
  3278. ...info.source === undefined ? {} : { source: info.source },
  3279. writable: info.writable,
  3280. }
  3281. return [ref, view] as const
  3282. }))
  3283. return ok(request, { credentials: Object.fromEntries(entries) })
  3284. },
  3285. async set(request) {
  3286. const credentials = ctx.get('credentials')
  3287. if (credentials === undefined) return err(request, credentialsAbsent())
  3288. const { ref, value } = request.payload
  3289. try {
  3290. await credentials.set(credentialRef(ref), value)
  3291. } catch (error: unknown) {
  3292. return err(request, {
  3293. code: 'credential-rejected',
  3294. message: error instanceof Error ? error.message : String(error),
  3295. details: { ref },
  3296. })
  3297. }
  3298. return ok(request, {})
  3299. },
  3300. async unset(request) {
  3301. const credentials = ctx.get('credentials')
  3302. if (credentials === undefined) return err(request, credentialsAbsent())
  3303. const { ref } = request.payload
  3304. try {
  3305. await credentials.unset(credentialRef(ref))
  3306. } catch (error: unknown) {
  3307. return err(request, {
  3308. code: 'credential-rejected',
  3309. message: error instanceof Error ? error.message : String(error),
  3310. details: { ref },
  3311. })
  3312. }
  3313. return ok(request, {})
  3314. },
  3315. },
  3316. llm: {
  3317. providers(request) {
  3318. const registered = ctx.llm.listProviders()
  3319. const active = new Set(registered.map(provider => provider.id))
  3320. const directory = ctx.llm.listConfigurableProviders()
  3321. const declared = new Set(directory.map(entry => entry.provider))
  3322. const views: ConfigurableProviderView[] = directory.map(entry => ({
  3323. provider: entry.provider,
  3324. displayName: entry.displayName,
  3325. settingsNs: entry.settingsNs,
  3326. settingsPath: [...entry.settingsPath],
  3327. active: active.has(entry.provider),
  3328. ...entry.declared === undefined ? {} : { declared: entry.declared },
  3329. }))
  3330. // Routes registered without a directory declaration still appear —
  3331. // they exist and serve models — just with no settings address. No
  3332. // adapter claimed them, so nothing can say whether they are shipped.
  3333. for (const provider of registered) {
  3334. if (declared.has(provider.id)) continue
  3335. views.push({
  3336. provider: provider.id,
  3337. displayName: provider.name,
  3338. settingsNs: '',
  3339. settingsPath: [],
  3340. active: true,
  3341. })
  3342. }
  3343. return Promise.resolve(ok(request, { providers: views }))
  3344. },
  3345. async models(request) {
  3346. return ok(request, await buildModelCatalog(ctx))
  3347. },
  3348. async discoverModels(request, signal) {
  3349. const { settingsNs, provider, baseURL, api, apiKey } = request.payload
  3350. try {
  3351. const models = await ctx.llm.discoverModels(settingsNs, {
  3352. ...provider === undefined ? {} : { provider },
  3353. ...baseURL === undefined ? {} : { baseURL },
  3354. ...api === undefined ? {} : { api },
  3355. ...apiKey === undefined ? {} : { apiKey },
  3356. ...signal === undefined ? {} : { signal },
  3357. })
  3358. return ok(request, { models })
  3359. } catch (error: unknown) {
  3360. // Every failure here is the user's next move, not a transport fault:
  3361. // a wrong endpoint, a rejected key, or a protocol with no listing all
  3362. // end at the same place — fill the models in by hand. The details
  3363. // repeat only what the caller already sent, never the credential.
  3364. return err(request, {
  3365. code: 'model-discovery-failed',
  3366. message: error instanceof Error ? error.message : String(error),
  3367. details: { settingsNs, ...baseURL === undefined ? {} : { baseURL } },
  3368. })
  3369. }
  3370. },
  3371. },
  3372. events: {
  3373. mux(_request, signal) {
  3374. const queue = new FrameQueue<RpcRequest<MuxFrame>>()
  3375. muxQueues.add(queue)
  3376. for (const session of ctx.sessions.list()) {
  3377. subscribeSession(queue, session)
  3378. }
  3379. for (const pending of pendingQuestions.values()) {
  3380. queue.push({
  3381. rpcId: pending.rpcId,
  3382. payload: {
  3383. type: 'question/requested', sessionId: pending.sessionId,
  3384. questions: pending.questions,
  3385. },
  3386. })
  3387. }
  3388. // Refresh recovery: still-pending approval questions replay with their
  3389. // stable rpcId so a reconnecting client can still answer them.
  3390. for (const pending of pendingApprovals.values()) queue.push(requestedFrame(pending))
  3391. // Queue snapshot baseline (pendingQuestions precedent): frames replayed
  3392. // in arrival order per session; a reconnecting client rebuilds its
  3393. // queue view from these alone.
  3394. for (const session of ctx.sessions.list()) {
  3395. const agent = ctx.agents.get(session.id)
  3396. if (agent?.session === session && agent.inbox.hasPending) {
  3397. queue.push(frame({ type: 'session/queue', sessionId: session.id, items: queueItems(agent) }))
  3398. }
  3399. }
  3400. // Background-task baseline. `ctx.agents.get` is the non-resuming read:
  3401. // a session with no live Agent owns no tasks, so it correctly sees only
  3402. // the unowned ones, and listing never revives a cold session. An empty
  3403. // set sends nothing — absence is how the client reads "no tasks".
  3404. const jobs = ctx.get('jobs')
  3405. if (jobs !== undefined) {
  3406. for (const session of ctx.sessions.list()) {
  3407. const views = jobViews(jobs.list(ctx.agents.get(session.id)))
  3408. if (views.length > 0) {
  3409. queue.push(frame({ type: 'session/jobs', sessionId: session.id, jobs: views }))
  3410. }
  3411. }
  3412. }
  3413. // Per-session open-call table for result-view pairing. Bounded by the
  3414. // per-turn call count: entries clear on turn/end; a table miss (stream
  3415. // opened mid-turn) backscans the session's in-memory events instead.
  3416. const openCalls = new Map<SessionId, Map<string, { name: string; args: unknown }>>()
  3417. const disposers = [
  3418. ctx.on('session/event', (session: Session, event: SessionEvent) => {
  3419. if (event.type === 'tool/call') {
  3420. const data = event.data as ToolCallData
  3421. try {
  3422. let table = openCalls.get(session.id)
  3423. if (table === undefined) openCalls.set(session.id, table = new Map<string, { name: string; args: unknown }>())
  3424. table.set(data.callId, { name: data.name, args: JSON.parse(data.arguments) })
  3425. } catch {
  3426. // Unparseable model arguments: leave the table unset; the result view soft-falls.
  3427. }
  3428. } else if (event.type === 'turn/end') {
  3429. openCalls.delete(session.id)
  3430. }
  3431. const view = viewFor(
  3432. ctx, event,
  3433. callId => openCalls.get(session.id)?.get(callId) ?? backscanArgs(session.events, callId),
  3434. ctx.agents.get(session.id),
  3435. )
  3436. queue.push(frame({ type: 'session/event', sessionId: session.id, event, ...view === undefined ? {} : { view } }))
  3437. }),
  3438. ctx.on('session/created', (session: Session) => {
  3439. subscribeSession(queue, session)
  3440. // The subscribe frame clears the client's task mirror, and a
  3441. // session born after the stream opened missed the baseline loop.
  3442. // Unowned tasks are visible to it from birth, so without this it
  3443. // would show none until the next registry change.
  3444. const views = jobs === undefined ? [] : jobViews(jobs.list(ctx.agents.get(session.id)))
  3445. if (views.length > 0) {
  3446. queue.push(frame({ type: 'session/jobs', sessionId: session.id, jobs: views }))
  3447. }
  3448. }),
  3449. ctx.on('session/disposed', (session: Session) => {
  3450. openCalls.delete(session.id)
  3451. }),
  3452. ...jobs === undefined ? [] : [jobs.onJobsChanged((owner) => {
  3453. if (owner !== undefined) {
  3454. // The exact owner instance the fence compares against, so the
  3455. // push stays correct even while that Agent's scope is tearing
  3456. // down and a lookup by id would already miss.
  3457. queue.push(frame({ type: 'session/jobs', sessionId: owner.id, jobs: jobViews(jobs.list(owner)) }))
  3458. return
  3459. }
  3460. // An unowned task is visible to every caller, so every subscribed
  3461. // session's set changed with it.
  3462. for (const session of ctx.sessions.list()) {
  3463. queue.push(frame({
  3464. type: 'session/jobs',
  3465. sessionId: session.id,
  3466. jobs: jobViews(jobs.list(ctx.agents.get(session.id))),
  3467. }))
  3468. }
  3469. })],
  3470. ]
  3471. return queue.iterate(signal, () => {
  3472. muxQueues.delete(queue)
  3473. for (const dispose of disposers) dispose()
  3474. })
  3475. },
  3476. host(_request, signal) {
  3477. const queue = new FrameQueue<RpcRequest<HostFrame>>()
  3478. const committedWorkspaces = ctx.workspaceRegistry.list()
  3479. const committedWorkspaceIds = new Set(
  3480. committedWorkspaces.map(workspace => String(workspace.id)),
  3481. )
  3482. let committedWorkspaceOrder = committedWorkspaces.map(workspace => workspace.id)
  3483. // Frame-dedup baseline, same posture as committedWorkspaceIds: the
  3484. // stream opens against the current set; workspace.list re-baselines
  3485. // reconnecting clients, so only later changes need frames.
  3486. let archivedSessionIds = ctx.workspaceRegistry.archivedSessionIds
  3487. const disposers = [
  3488. ctx.on('session/created', (session: Session) => {
  3489. queue.push(frame({
  3490. type: 'host/session-added',
  3491. sessionId: session.id,
  3492. // Derived at frame time like summarize(); a just-created session
  3493. // has run no turn yet, so this is constantly true in practice.
  3494. blank: sessionBlank(session),
  3495. // Including cwd lets the client group the new session without refreshing the list.
  3496. ...sessionListFields(session.header, session.events),
  3497. }))
  3498. }),
  3499. ctx.on('session/disposed', (session: Session) => {
  3500. queue.push(frame({ type: 'host/session-removed', sessionId: session.id }))
  3501. }),
  3502. ctx.on('agent/status', ({ agent, status }: { agent: Agent; status: AgentStatus }) => {
  3503. queue.push(frame({ type: 'host/session-status', sessionId: agent.id, running: status === 'running' }))
  3504. }),
  3505. ctx.on('agent/error', ({ agent, error }: { agent: Agent; error: unknown }) => {
  3506. queue.push(frame({ type: 'host/agent-error', sessionId: agent.id, message: errorChain(error) }))
  3507. }),
  3508. ctx.on('domain/changed', (change) => {
  3509. if (change.domain !== 'workspace') return
  3510. if (change.table === '') {
  3511. if (change.operation !== 'put') return
  3512. const state = workspaceDomainState.parse(change.value)
  3513. const orderChanged = state.workspaceIds.length === committedWorkspaceOrder.length
  3514. && state.workspaceIds.every(workspaceId => committedWorkspaceIds.has(String(workspaceId)))
  3515. && state.workspaceIds.some((workspaceId, index) => workspaceId !== committedWorkspaceOrder[index])
  3516. for (const workspaceId of state.workspaceIds) {
  3517. if (committedWorkspaceIds.has(workspaceId)) continue
  3518. const workspace = ctx.workspaceRegistry.get(workspaceId)
  3519. if (workspace === undefined) {
  3520. throw new Error(`committed workspace registry references missing workspace "${workspaceId}"`)
  3521. }
  3522. committedWorkspaceIds.add(workspaceId)
  3523. queue.push(frame({ type: 'host/workspace-changed', workspace: workspaceView(workspace) }))
  3524. }
  3525. committedWorkspaceOrder = [...state.workspaceIds]
  3526. if (orderChanged) {
  3527. queue.push(frame({
  3528. type: 'host/workspace-order-changed',
  3529. workspaceIds: [...state.workspaceIds],
  3530. }))
  3531. }
  3532. if (state.archivedSessionIds.length !== archivedSessionIds.length
  3533. || state.archivedSessionIds.some((id, index) => id !== archivedSessionIds[index])) {
  3534. archivedSessionIds = state.archivedSessionIds
  3535. queue.push(frame({
  3536. type: 'host/archived-sessions-changed',
  3537. archivedSessionIds: [...state.archivedSessionIds],
  3538. }))
  3539. }
  3540. return
  3541. }
  3542. if (change.table !== 'workspaces') return
  3543. if (change.operation === 'deleted') {
  3544. if (!committedWorkspaceIds.delete(change.key)) return
  3545. queue.push(frame({
  3546. type: 'host/workspace-removed',
  3547. workspaceId: change.key as WorkspaceId,
  3548. }))
  3549. return
  3550. }
  3551. if (!committedWorkspaceIds.has(change.key)) return
  3552. // Existing-entity table writes are complete attach/touch commits.
  3553. // A new entity's first put waits for the global registry write above.
  3554. queue.push(frame({
  3555. type: 'host/workspace-changed',
  3556. workspace: changedWorkspaceView(change.key, change.value),
  3557. }))
  3558. }),
  3559. // Allowlisted host events ride one verbatim wrapper frame each. The
  3560. // allowlist is api-remotes', and `ctx.remote.$on` is the consumer
  3561. // face; nothing here projects, redacts, or renames.
  3562. ...API_REMOTE_FORWARDED_EVENTS.map(name => ctx.on(
  3563. name,
  3564. // The allowlist's shape assertion proves each name is a real,
  3565. // non-scoped, void-returning event, so the rest-parameter handler
  3566. // satisfies every member of the union `on` accepts here;
  3567. // assertJsonArgs proves the payload is JSON-safe before it queues.
  3568. ((...args: unknown[]) => {
  3569. queue.push(frame({
  3570. type: 'host/remote-event',
  3571. event: name,
  3572. args: assertJsonArgs(name, args),
  3573. }))
  3574. }),
  3575. )),
  3576. ]
  3577. return queue.iterate(signal, () => { for (const dispose of disposers) dispose() })
  3578. },
  3579. },
  3580. downloads: {
  3581. async sessionLog(request, signal) {
  3582. // Clean error path first: missing services answer 500 and a missing
  3583. // root artifact 404 before any zip byte is produced. The root content
  3584. // read here is reused as the first zip entry, so nothing is read twice.
  3585. const deps = sessionLogExportDeps(ctx)
  3586. if (deps.sessionQuery === undefined || deps.sessionPersistence === undefined || deps.attachments === undefined) {
  3587. return new Response(
  3588. 'session log export is unavailable: missing session-query, session-persistence, or attachments service',
  3589. { status: 500 },
  3590. )
  3591. }
  3592. if (!deps.sessionPersistence.supportsRawArtifacts) {
  3593. return new Response(
  3594. 'session log export is unavailable: the persistence backend does not expose per-session raw artifacts',
  3595. { status: 501 },
  3596. )
  3597. }
  3598. const ready: SessionLogExportReady = {
  3599. sessionQuery: deps.sessionQuery,
  3600. sessionPersistence: deps.sessionPersistence,
  3601. attachments: deps.attachments,
  3602. sessions: deps.sessions,
  3603. }
  3604. let root: SessionRawArtifact | undefined
  3605. try {
  3606. await flushLiveSessionLog(deps, request.sessionId, signal)
  3607. root = await deps.sessionPersistence.readRaw(request.sessionId, signal)
  3608. signal.throwIfAborted()
  3609. } catch {
  3610. signal.throwIfAborted()
  3611. // Root preparation failure: answer 500 without echoing the error,
  3612. // which may carry absolute host paths into the browser error bar.
  3613. return new Response('session log export failed to prepare the stored artifact', { status: 500 })
  3614. }
  3615. if (root === undefined) {
  3616. return new Response('session not found', { status: 404 })
  3617. }
  3618. return new Response(
  3619. streamSessionLogZip(
  3620. ready,
  3621. root,
  3622. request.sessionId,
  3623. request.includeDescendants === true,
  3624. sessionExportCompressionLevel,
  3625. signal,
  3626. ),
  3627. {
  3628. headers: {
  3629. 'content-type': 'application/zip',
  3630. 'content-disposition': `attachment; filename="${sessionLogZipFilename(request.sessionId)}"`,
  3631. },
  3632. },
  3633. )
  3634. },
  3635. },
  3636. respond(message: ClientResponse): Promise<RpcReceipt> {
  3637. // Route by the echoed rpcId (the wire correlation): approvals first,
  3638. // then questions — the two registries share one id space of UUIDs.
  3639. const approval = pendingApprovals.get(message.rpcId)
  3640. if (approval !== undefined) {
  3641. if (!message.result.ok) return Promise.resolve({ accepted: false, reason: 'bad-response' })
  3642. const parsed = approvalResponsePayloadSchema.safeParse(message.result.value)
  3643. // The payload's audit correlation must match the entry the rpcId routed
  3644. // to — a mismatched answer is malformed, not merely late.
  3645. if (!parsed.success || parsed.data.approvalId !== approval.approvalId || parsed.data.sessionId !== approval.sessionId) {
  3646. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  3647. }
  3648. approval.resolve(parsed.data.outcome)
  3649. return Promise.resolve({ accepted: true })
  3650. }
  3651. const pending = pendingQuestions.get(message.rpcId)
  3652. if (pending === undefined) return Promise.resolve({ accepted: false, reason: 'not-pending' })
  3653. if (!message.result.ok) {
  3654. if (message.result.error.code !== 'cancelled') {
  3655. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  3656. }
  3657. claimQuestion(pending, 'cancelled')
  3658. pending.reject(new UserQuestionError(
  3659. 'the user cancelled ask_user_question', 'ASK_CANCELLED'))
  3660. return Promise.resolve({ accepted: true })
  3661. }
  3662. const parsed = questionResponsePayloadSchema.safeParse(message.result.value)
  3663. if (!parsed.success) {
  3664. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  3665. }
  3666. const payload: QuestionResponsePayload = {
  3667. sessionId: parsed.data.sessionId,
  3668. answer: {
  3669. answers: parsed.data.answer.answers.map(answer => ({
  3670. id: answer.id,
  3671. selected: answer.selected,
  3672. ...(answer.custom === undefined ? {} : { custom: answer.custom }),
  3673. })),
  3674. },
  3675. }
  3676. if (!matchesQuestions(payload, pending)) {
  3677. return Promise.resolve({ accepted: false, reason: 'bad-response' })
  3678. }
  3679. claimQuestion(pending, 'answered')
  3680. pending.resolve(payload.answer)
  3681. return Promise.resolve({ accepted: true })
  3682. },
  3683. }
  3684. }