返回源码地图

packages/api/session-controller/src/agent.ts

main snapshot · da00f7f5358f · 正文引用章节 14 / 15;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/** Agent activation, composition, and model-selection policy owned by API Session. */
2
3import { mkdir } from 'node:fs/promises'
4import type { Context } from '@deepseek-ai/cordis'
5import { installModelSelection } from '@deepseek-ai/dsh-agent'
6import type {
7 Agent, AgentOptions, AgentSetup, ModelSelection as AgentModelSelection, ModelSelectionRef,
8} from '@deepseek-ai/dsh-agent'
9import type {} from '@deepseek-ai/dsh-agent-default-model'
10import type {} from '@deepseek-ai/dsh-agent-preset-registry'
11import { ReasoningEffortId } from '@deepseek-ai/dsh-llm'
12import type { Session, SessionId } from '@deepseek-ai/dsh-session'
13import type { SessionInspection } from '@deepseek-ai/dsh-session-persistence'
14import { SessionQueryError, type SessionObservation } from '@deepseek-ai/dsh-session-query'
15import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
16import type {} from '@deepseek-ai/dsh-typert-registry'
17import type { ModelSelection } from './types.ts'
18
19/** Cold Session identity absent from persistence. */
20export class ApiSessionNotFound extends Error {}
21
22/** Session identity whose lifecycle belongs to subagent routing. */
23export class ApiSessionSubagentOwnership extends Error {
24 /** @param sessionId - identity reserved to subagent routing. */
25 constructor(readonly sessionId: SessionId) {
26 super(`session "${sessionId}" is a subagent session; use subagent delivery`)
27 }
28}
29
30/** Explicit-id creation attempted to adopt a Session under another cwd. */
31export class ApiSessionCwdConflict extends Error {
32 constructor(
33 readonly sessionId: SessionId,
34 readonly requestedCwd: string,
35 readonly existingCwd: string | undefined,
36 ) {
37 super(
38 existingCwd === undefined
39 ? `session "${sessionId}" records no cwd and cannot be adopted for "${requestedCwd}"`
40 : `session "${sessionId}" belongs to "${existingCwd}", not "${requestedCwd}"`,
41 )
42 }
43}
44
45/** Explicit-id creation attempted to adopt a Session under another preset. */
46export class ApiSessionPresetConflict extends Error {
47 constructor(
48 readonly sessionId: SessionId,
49 readonly requestedPreset: string,
50 readonly existingPreset: string | undefined,
51 ) {
52 super(
53 existingPreset === undefined
54 ? `session "${sessionId}" records no agent preset and cannot be adopted under "${requestedPreset}"`
55 : `session "${sessionId}" runs agent preset "${existingPreset}", not "${requestedPreset}"`,
56 )
57 }
58}
59
60/** Failures produced while resolving one ordinary Session identity to its live Agent. */
61export type ApiSessionAgentError = RemoteError<'session/not-found' | 'session/agent-busy' | 'session/writer-held' | 'gateway/internal'>
62
63/** Result of resolving one ordinary Session identity to its live Agent. */
64export type ApiSessionAgentResult =
65 | { readonly agent: Agent }
66 | { readonly error: ApiSessionAgentError }
67
68type InstalledSelection = ModelSelectionRef & {
69 current: AgentModelSelection
70 consume(provider: string, model: string, reasoningEffort: string | undefined): boolean
71}
72
73/**
74 * Test whether generic Session routing must leave an identity to subagent routing.
75 * @param ctx - Host context carrying the Agent ownership registry.
76 * @param session - attached or live Session whose ownership is tested.
77 * @param agent - live Agent when one exists for the Session.
78 * @returns whether subagent routing owns the Session identity.
79 */
80export function hasApiSessionSubagentOwner(
81 ctx: Context,
82 session: Pick<Session, 'header'>,
83 agent: Agent | undefined,
84): boolean {
85 if (session.header.origin === 'subagent') return true
86 const parentId = session.header.parentSession
87 if (parentId === undefined || agent === undefined) return false
88 const parent = ctx.agents.get(parentId)
89 return parent !== undefined && ctx.agents.isOwnedBy(agent.id, parent)
90}
91
92/**
93 * Build the stable caller-facing subagent ownership rejection.
94 * @param sessionId - Session identity owned by subagent routing.
95 * @returns a stable Session-domain failure.
96 */
97export function apiSessionSubagentOwnershipError(sessionId: SessionId): ApiSessionAgentError {
98 return new RemoteError(
99 'session/agent-busy',
100 `session "${sessionId}" is owned by subagent routing`,
101 { reason: 'use subagent delivery for this child session' },
102 )
103}
104
105/**
106 * Inspect one cold Session without repairing, resuming, or publishing it.
107 * @param ctx - Host context carrying Session persistence.
108 * @param sessionId - durable Session identity.
109 * @param signal - optional cancellation for persistence reads.
110 * @returns the persisted header and complete event prefix.
111 */
112export async function inspectApiSession(
113 ctx: Context,
114 sessionId: SessionId,
115 signal?: AbortSignal,
116): Promise<SessionInspection> {
117 try {
118 using observation = await ctx.sessionQuery.observeSession(sessionId, {
119 ...(signal === undefined ? {} : { signal }),
120 projectionMode: 'none',
121 })
122 if (observation.header.cwd === undefined) {
123 throw new ApiSessionNotFound(`session "${sessionId}" not found`)
124 }
125 return {
126 meta: observation.header,
127 inheritedEventCount: observation.inheritedEventCount,
128 events: [...observation.events],
129 }
130 } catch (error: unknown) {
131 if (error instanceof SessionQueryError
132 && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') {
133 throw new ApiSessionNotFound(`session "${sessionId}" not found`)
134 }
135 throw error
136 }
137}
138
139/** Owns every operation that may create, resume, or configure a Web Agent. */
140export class ApiSessionAgentController {
141 private readonly resumes = new Map<SessionId, Promise<Agent>>()
142 private readonly creations = new Map<SessionId, Promise<Agent>>()
143 private readonly selections = new WeakMap<Agent, InstalledSelection>()
144 private readonly imageAdmissionChains = new WeakMap<Agent, Promise<void>>()
145
146 /** @param ctx - Host context carrying Agent, model, persistence, and Typert services. */
147 constructor(private readonly ctx: Context) {
148 ctx.typert.lookups.configure('agent', async (sessionId: SessionId) => {
149 const found = await this.resolveAgent(sessionId)
150 if ('error' in found) throw found.error
151 return found.agent
152 })
153 ctx.typert.lookups.configure('session', async (sessionId: SessionId) => {
154 const found = await this.resolveAgent(sessionId)
155 if ('error' in found) throw found.error
156 return found.agent.session
157 })
158 ctx.typert.contexts.configureHost('agent', async (sessionId: SessionId) => {
159 const found = await this.resolveAgent(sessionId)
160 if ('error' in found) throw found.error
161 return found.agent.ctx
162 })
163 }
164
165 /**
166 * Resolve or resume one ordinary Session, deduplicating concurrent resumes.
167 * @param sessionId - ordinary Session identity.
168 * @returns the live Agent or a stable Session-domain failure.
169 */
170 async resolveAgent(sessionId: SessionId): Promise<ApiSessionAgentResult> {
171 return this.resolve(sessionId)
172 }
173
174 /**
175 * Resolve one ordinary Session from an already-retained exact observation.
176 * @param observation - Host-owned observation whose preparation stays pinned through setup.
177 * @returns the live Agent or a stable Session-domain failure.
178 */
179 async resolveObservedAgent(observation: SessionObservation): Promise<ApiSessionAgentResult> {
180 return this.resolve(observation.header.id, observation)
181 }
182
183 private async resolve(
184 sessionId: SessionId,
185 observation?: SessionObservation,
186 ): Promise<ApiSessionAgentResult> {
187 const live = this.liveAgent(sessionId)
188 if (live !== undefined) return live
189 const attached = this.ctx.sessions.get(sessionId)
190 if (attached !== undefined && hasApiSessionSubagentOwner(this.ctx, attached, undefined)) {
191 return { error: apiSessionSubagentOwnershipError(sessionId) }
192 }
193
194 let resume = this.resumes.get(sessionId)
195 if (resume === undefined) {
196 resume = this.resume(sessionId, observation).finally(() => { this.resumes.delete(sessionId) })
197 this.resumes.set(sessionId, resume)
198 }
199 try {
200 const agent = await resume
201 // A shared resume can publish an identity that subagent routing adopts
202 // before every waiter observes it; apply the live ownership policy again.
203 const published = this.liveAgent(sessionId)
204 return published ?? { agent }
205 } catch (error: unknown) {
206 if (error instanceof ApiSessionNotFound) {
207 return { error: new RemoteError('session/not-found', error.message, { sessionId }) }
208 }
209 if (error instanceof ApiSessionSubagentOwnership) {
210 return { error: apiSessionSubagentOwnershipError(error.sessionId) }
211 }
212 const raced = this.liveAgent(sessionId)
213 if (raced !== undefined) return raced
214 const racedSession = this.ctx.sessions.get(sessionId)
215 if (racedSession !== undefined && hasApiSessionSubagentOwner(this.ctx, racedSession, undefined)) {
216 return { error: apiSessionSubagentOwnershipError(sessionId) }
217 }
218 if (error instanceof Error && error.name === 'SessionAlreadyOwnedError') {
219 return { error: new RemoteError('session/writer-held', error.message, { sessionId }) }
220 }
221 return {
222 error: new RemoteError(
223 'gateway/internal',
224 `resume failed for session "${sessionId}": ${String(error)}`,
225 {},
226 ),
227 }
228 }
229 }
230
231 /**
232 * Resolve one requested identity, creating or resuming it once.
233 * @param sessionId - requested Session identity.
234 * @param cwd - directory the Session must own.
235 * @param checkPersistedIdentity - whether to inspect a cold identity before creation.
236 * @param presetId - optional Agent preset the Session must own.
237 * @returns the matching live ordinary Agent.
238 */
239 async ensureSession(
240 sessionId: SessionId,
241 cwd: string,
242 checkPersistedIdentity: boolean,
243 presetId?: string,
244 ): Promise<Agent> {
245 let creation = this.creations.get(sessionId)
246 if (creation === undefined) {
247 creation = this.createOrAdopt(sessionId, cwd, checkPersistedIdentity, presetId)
248 .catch((error: unknown) => {
249 const live = this.ctx.agents.get(sessionId)
250 if (live !== undefined) {
251 if (hasApiSessionSubagentOwner(this.ctx, live.session, live)) {
252 throw new ApiSessionSubagentOwnership(sessionId)
253 }
254 return live
255 }
256 const attached = this.ctx.sessions.get(sessionId)
257 if (attached !== undefined && hasApiSessionSubagentOwner(this.ctx, attached, undefined)) {
258 throw new ApiSessionSubagentOwnership(sessionId)
259 }
260 throw error
261 })
262 .finally(() => { this.creations.delete(sessionId) })
263 this.creations.set(sessionId, creation)
264 }
265 const agent = await creation
266 if (hasApiSessionSubagentOwner(this.ctx, agent.session, agent)) {
267 throw new ApiSessionSubagentOwnership(sessionId)
268 }
269 if (presetId !== undefined) {
270 this.assertPresetUnchanged(sessionId, presetId, this.presetForSession(agent.session))
271 }
272 if (agent.session.header.cwd !== cwd) {
273 throw new ApiSessionCwdConflict(sessionId, cwd, agent.session.header.cwd)
274 }
275 return agent
276 }
277
278 /**
279 * Install or return the Session-local model selection used by prompt assembly.
280 * @param agent - live Agent that owns the selection.
281 * @returns the installed mutable selection reference.
282 */
283 selectionFor(agent: Agent): InstalledSelection {
284 const installed = this.selections.get(agent)
285 if (installed !== undefined) return installed
286 const projectionState = this.ctx.sessionProjections.stateOf(agent.session, 'modelSelection')
287 if (projectionState === undefined) {
288 throw new Error('api-session: required modelSelection projection is not registered')
289 }
290 let picked = projectionState.pending === null
291 ? undefined
292 : agentModelSelection(projectionState.pending)
293 const defaultModel = this.ctx.agentDefaultModel
294 const selection: InstalledSelection = {
295 get current(): AgentModelSelection {
296 if (picked !== undefined) return picked
297 const loggedHeader = agent.session.requestHeader()
298 if (loggedHeader === undefined) return defaultModel.currentSelection()
299 const logged = loggedHeader.config
300 return {
301 provider: logged.provider,
302 model: logged.model,
303 // An effort the adapter defaulted is not a conversation choice: restoring
304 // it as one would make an unchanged default read as a request change.
305 ...(logged.reasoningEffort === undefined
306 || loggedHeader.adapterDefaults?.reasoningEffort === true
307 ? {}
308 : { reasoningEffort: logged.reasoningEffort }),
309 }
310 },
311 set current(next: AgentModelSelection) {
312 picked = next
313 },
314 consume(provider: string, model: string, reasoningEffort: string | undefined): boolean {
315 if (picked?.provider !== provider
316 || picked.model !== model
317 || picked.reasoningEffort !== reasoningEffort) return false
318 picked = undefined
319 return true
320 },
321 assembled: undefined,
322 }
323 installModelSelection(agent.ctx, selection)
324 this.selections.set(agent, selection)
325 return selection
326 }
327
328 /**
329 * Commit and cache one validated selection for the next prompt assembly.
330 * @param agent - live Agent that owns the selection.
331 * @param selection - validated selection to record and apply.
332 */
333 selectForNextRequest(agent: Agent, selection: AgentModelSelection): void {
334 agent.session.append('model/selection', selection)
335 this.selectionFor(agent).current = selection
336 }
337
338 /**
339 * Let a matching durable request header retire the execution cache.
340 * @param agent - live Agent whose request was recorded.
341 * @param provider - provider route used by the request.
342 * @param model - provider-owned model used by the request.
343 * @param reasoningEffort - adapter-owned effort used by the request.
344 * @returns whether the pending selection was consumed.
345 */
346 consumeSelection(
347 agent: Agent,
348 provider: string,
349 model: string,
350 reasoningEffort: string | undefined,
351 ): boolean {
352 return this.selections.get(agent)?.consume(provider, model, reasoningEffort) ?? false
353 }
354
355 /**
356 * Read the current Agent preset from the Session projection.
357 * @param session - live Session whose projection state is available.
358 * @returns the current preset, or undefined when the capability is absent.
359 */
360 presetForSession(session: Session): string | undefined {
361 return this.ctx.sessionProjections.stateOf(session, 'agentPreset') ?? undefined
362 }
363
364 /**
365 * Serialize image admission and model selection for one Agent.
366 * @param agent - live Agent that owns the serialization chain.
367 * @param operation - asynchronous operation admitted after prior work settles.
368 * @returns the operation result or rejection.
369 */
370 serializeImageAdmission<Value>(agent: Agent, operation: () => Promise<Value>): Promise<Value> {
371 const result = (this.imageAdmissionChains.get(agent) ?? Promise.resolve()).then(operation)
372 this.imageAdmissionChains.set(agent, result.then(() => undefined, () => undefined))
373 return result
374 }
375
376 /**
377 * Resolve the preset id and pre-publication Agent setup for a create or resume.
378 * @param presetId - requested preset or the configured default when omitted.
379 * @returns the resolved preset identity and Agent setup callback.
380 */
381 async composeAgent(presetId: string | undefined): Promise<{
382 readonly agentPreset?: string
383 readonly setup: AgentSetup
384 }> {
385 const presets = this.ctx.get('agentPresets')
386 if (presets === undefined) {
387 return { setup: (_agentCtx, agent) => { this.installSelection(agent) } }
388 }
389 const resolvedId = (await presets.resolve(presetId)).id
390 return {
391 agentPreset: resolvedId,
392 setup: async (agentCtx, agent) => {
393 this.installSelection(agent)
394 await presets.mount(agentCtx, resolvedId)
395 },
396 }
397 }
398
399 private liveAgent(sessionId: SessionId): ApiSessionAgentResult | undefined {
400 const agent = this.ctx.agents.get(sessionId)
401 if (agent === undefined) return undefined
402 return hasApiSessionSubagentOwner(this.ctx, agent.session, agent)
403 ? { error: apiSessionSubagentOwnershipError(sessionId) }
404 : { agent }
405 }
406
407 private async resume(sessionId: SessionId, supplied?: SessionObservation): Promise<Agent> {
408 if (supplied !== undefined) return this.resumeObserved(sessionId, supplied)
409 try {
410 using observation = await this.ctx.sessionQuery.observeSession(sessionId)
411 return await this.resumeObserved(sessionId, observation)
412 } catch (error: unknown) {
413 if (error instanceof SessionQueryError
414 && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') {
415 throw new ApiSessionNotFound(`session "${sessionId}" not found`)
416 }
417 throw error
418 }
419 }
420
421 private async resumeObserved(
422 sessionId: SessionId,
423 observation: SessionObservation,
424 ): Promise<Agent> {
425 if (observation.header.id !== sessionId || observation.header.cwd === undefined) {
426 throw new ApiSessionNotFound(`session "${sessionId}" not found`)
427 }
428 if (hasApiSessionSubagentOwner(this.ctx, { header: observation.header }, undefined)) {
429 throw new ApiSessionSubagentOwnership(sessionId)
430 }
431 const composition = await this.composeAgent(this.presetForObservation(observation))
432 const published = this.ctx.sessions.get(sessionId)
433 const live = this.ctx.agents.get(sessionId)
434 if (published !== undefined && hasApiSessionSubagentOwner(this.ctx, published, live)) {
435 throw new ApiSessionSubagentOwnership(sessionId)
436 }
437 return (await this.ctx.agents.resume({
438 resumeSessionId: sessionId,
439 agentOptions: this.agentOptions(),
440 setup: composition.setup,
441 })).agent
442 }
443
444 private async createOrAdopt(
445 sessionId: SessionId,
446 cwd: string,
447 checkPersistedIdentity: boolean,
448 presetId: string | undefined,
449 ): Promise<Agent> {
450 const attached = this.ctx.sessions.get(sessionId)
451 const live = this.ctx.agents.get(sessionId)
452 if (attached !== undefined && hasApiSessionSubagentOwner(this.ctx, attached, live)) {
453 throw new ApiSessionSubagentOwnership(sessionId)
454 }
455 if (live !== undefined) return live
456
457 if (checkPersistedIdentity) {
458 try {
459 using observation = await this.ctx.sessionQuery.observeSession(sessionId)
460 if (hasApiSessionSubagentOwner(this.ctx, { header: observation.header }, undefined)) {
461 throw new ApiSessionSubagentOwnership(sessionId)
462 }
463 if (observation.header.cwd !== cwd) {
464 throw new ApiSessionCwdConflict(sessionId, cwd, observation.header.cwd)
465 }
466 const storedPreset = this.presetForObservation(observation)
467 this.assertPresetUnchanged(sessionId, presetId, storedPreset)
468 const composition = await this.composeAgent(storedPreset)
469 return (await this.ctx.agents.resume({
470 resumeSessionId: sessionId,
471 agentOptions: this.agentOptions(),
472 setup: composition.setup,
473 })).agent
474 } catch (error: unknown) {
475 if (!(error instanceof SessionQueryError)
476 || error.code !== 'SESSION_QUERY_SESSION_NOT_FOUND') throw error
477 }
478 }
479
480 try {
481 await mkdir(cwd, { recursive: true })
482 } catch (error: unknown) {
483 throw new Error(`failed to ensure project directory "${cwd}": ${String(error)}`, { cause: error })
484 }
485 const composition = await this.composeAgent(presetId)
486 return (await this.ctx.agents.create({
487 sessionId,
488 agentOptions: this.agentOptions(),
489 meta: {
490 cwd,
491 ...(composition.agentPreset === undefined ? {} : { agentPreset: composition.agentPreset }),
492 },
493 setup: composition.setup,
494 })).agent
495 }
496
497 private agentOptions(): AgentOptions {
498 const { provider, model } = this.ctx.agentDefaultModel.currentSelection()
499 return { provider, model }
500 }
501
502 private installSelection(agent: Agent): void {
503 this.selectionFor(agent)
504 }
505
506 /**
507 * Read the current Agent preset from an all-projections observation.
508 * @param observation - exact Session observation carrying its projection snapshot.
509 * @returns the current preset, or undefined when the capability is absent.
510 */
511 presetForObservation(observation: SessionObservation): string | undefined {
512 if (observation.projections === undefined) {
513 throw new Error('api-session: Agent activation requires a projected Session observation')
514 }
515 return observation.projections.values.agentPreset ?? undefined
516 }
517
518 private assertPresetUnchanged(
519 sessionId: SessionId,
520 requested: string | undefined,
521 existing: string | undefined,
522 ): void {
523 if (requested === undefined || requested === existing) return
524 throw new ApiSessionPresetConflict(sessionId, requested, existing)
525 }
526}
527
528function agentModelSelection(selection: ModelSelection): AgentModelSelection {
529 return {
530 provider: selection.provider,
531 model: selection.model,
532 ...(selection.reasoningEffort === undefined
533 ? {}
534 : { reasoningEffort: ReasoningEffortId(selection.reasoningEffort) }),
535 }
536}