返回源码地图

packages/acp/acp/src/session.ts

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

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

1/** One standard ACP session's Agent, configuration, prompt, update, and teardown lifecycle. */
2
3import type { Context } from '@deepseek-ai/cordis'
4import {
5 RequestError,
6 type McpServer,
7 type PromptRequest,
8 type PromptResponse,
9 type SessionConfigOption,
10 type SessionNotification,
11 type StopReason,
12} from '@agentclientprotocol/sdk'
13import type { Agent, AgentHandle, AgentOptions, ModelSelection } from '@deepseek-ai/dsh-agent'
14import { createUserMessage, errorChain, type UserMessage } from '@deepseek-ai/dsh-llm'
15import { type Session, type SessionEvent, type SessionId, type TurnEndReason } from '@deepseek-ai/dsh-session'
16import { AcpContentError, admitAcpPrompt } from './content.ts'
17import { turnEndToStopReason } from './codec.ts'
18import { mountAcpMcpServers } from './mcp.ts'
19import { AcpModelControl } from './model-control.ts'
20import { assistantUpdates, toolCallUpdate, toolResultUpdate } from './updates.ts'
21
22/** The continuable-subagent teardown used without depending on the subagent package. */
23interface ContinuableDrain {
24 /** Dispose continuable descendants below exact host-owned parents child-first. */
25 drainContinuableDescendants(parents: readonly Agent[]): Promise<void>
26}
27
28/** Inputs shared by fresh and resumed ACP session construction. */
29interface AcpSessionBuildOptions {
30 cwd: string
31 mcpServers: readonly McpServer[]
32 agentOptions: AgentOptions
33 fallbackSelection: ModelSelection | undefined
34 signal: AbortSignal
35 notify: (notification: SessionNotification) => Promise<void>
36}
37
38/** Fresh ACP session construction inputs. */
39export interface CreateAcpSessionOptions extends AcpSessionBuildOptions {
40 sessionId: SessionId
41}
42
43/** Persisted ACP session construction inputs. */
44export interface ResumeAcpSessionOptions extends AcpSessionBuildOptions {
45 sessionId: SessionId
46}
47
48interface InflightPrompt {
49 resolve: (reason: StopReason) => void
50 reject: (error: Error) => void
51 messageId: string | undefined
52 messageQueued: boolean
53 turn: number | undefined
54 endReason: TurnEndReason | undefined
55 admissionDone: Promise<void>
56 finishAdmission: () => void
57 admissionController: AbortController
58 cancelRequested: boolean
59 settlementStarted: boolean
60 outputError: Error | undefined
61 agentError: Error | undefined
62}
63
64/** Standard invalid-parameter failure with protocol-safe detail. */
65function invalidParams(detail: string): RequestError {
66 return RequestError.invalidParams(undefined, detail)
67}
68
69/** Standard internal failure with protocol-safe detail. */
70function internalError(detail: string): RequestError {
71 return RequestError.internalError(undefined, detail)
72}
73
74/** Restore the latest logged route before falling back to deployment config. */
75function selectionFor(
76 logged: {
77 config: { provider: string; model: string; reasoningEffort?: ModelSelection['reasoningEffort'] }
78 adapterDefaults?: { reasoningEffort?: boolean }
79 } | undefined,
80 fallback: ModelSelection | undefined,
81): ModelSelection | undefined {
82 return logged === undefined
83 ? fallback
84 : {
85 provider: logged.config.provider,
86 model: logged.config.model,
87 ...logged.config.reasoningEffort === undefined || logged.adapterDefaults?.reasoningEffort === true
88 ? {}
89 : { reasoningEffort: logged.config.reasoningEffort },
90 }
91}
92
93/**
94 * Per-session ACP module. It owns the unpublished Agent composition, selected
95 * route, one-prompt admission slot, ordered standard updates, and memoized
96 * quiescent teardown.
97 */
98export class AcpSession {
99 /** The exact top-level Agent owned by this ACP session. */
100 readonly agent: Agent
101 private readonly modelControl: AcpModelControl
102 private outputTail = Promise.resolve()
103 private inflight: InflightPrompt | undefined
104 private closing: Promise<void> | undefined
105 private readonly pendingSelections = new Map<string, ModelSelection>()
106
107 private constructor(
108 private readonly ctx: Context,
109 handle: AgentHandle,
110 modelControl: AcpModelControl,
111 private readonly notify: (notification: SessionNotification) => Promise<void>,
112 ) {
113 this.agent = handle.agent
114 this.modelControl = modelControl
115 this.disposeAgent = () => handle.dispose()
116 }
117
118 private readonly disposeAgent: () => Promise<void>
119
120 /**
121 * Compose a fresh Agent and all requested MCP clients before publication.
122 * @param ctx - ACP plugin context with Agent, LLM, and persistence services.
123 * @param options - fresh session identity, workspace, route, MCP, and notifier.
124 * @returns the fully composed per-session module.
125 */
126 static async create(ctx: Context, options: CreateAcpSessionOptions): Promise<AcpSession> {
127 const modelControl = new AcpModelControl(ctx.llm, options.fallbackSelection)
128 const handle = await ctx.agents.create({
129 sessionId: options.sessionId,
130 meta: { cwd: options.cwd },
131 agentOptions: options.agentOptions,
132 signal: options.signal,
133 setup: async (agentCtx) => {
134 modelControl.install(agentCtx)
135 await mountAcpMcpServers(agentCtx, options.mcpServers, options.cwd)
136 },
137 })
138 return new AcpSession(ctx, handle, modelControl, options.notify)
139 }
140
141 /**
142 * Restore a persisted Agent and compose the request's fresh MCP connections.
143 * @param ctx - ACP plugin context with Agent, LLM, and persistence services.
144 * @param options - persisted identity, workspace, fallback route, MCP, and notifier.
145 * @returns the restored per-session module.
146 */
147 static async resume(ctx: Context, options: ResumeAcpSessionOptions): Promise<AcpSession> {
148 let modelControl: AcpModelControl | undefined
149 const handle = await ctx.agents.resume({
150 resumeSessionId: options.sessionId,
151 agentOptions: options.agentOptions,
152 signal: options.signal,
153 setup: async (agentCtx, agent) => {
154 modelControl = new AcpModelControl(
155 ctx.llm,
156 selectionFor(agent.session.requestHeader(), options.fallbackSelection),
157 )
158 modelControl.install(agentCtx)
159 await mountAcpMcpServers(agentCtx, options.mcpServers, options.cwd)
160 },
161 })
162 /* v8 ignore start -- a fulfilled Agent resume necessarily ran setup to completion. */
163 if (modelControl === undefined) {
164 await handle.dispose()
165 throw internalError('session/resume did not compose model selection')
166 }
167 /* v8 ignore stop */
168 return new AcpSession(ctx, handle, modelControl, options.notify)
169 }
170
171 /**
172 * Whether this module owns an exact Agent reference.
173 * @param agent - Agent observed on a scoped runtime event.
174 * @returns true only for this session's owned Agent.
175 */
176 owns(agent: Agent): boolean {
177 return this.agent === agent
178 }
179
180 /**
181 * Whether this module owns an exact Session reference.
182 * @param session - Session observed on a durable event.
183 * @returns true only for this session's owned Session.
184 */
185 ownsSession(session: Session): boolean {
186 return this.agent.session === session
187 }
188
189 /**
190 * Return the complete standard model configuration state.
191 * @param signal - optional request cancellation.
192 * @returns provider-grouped model and exact-model reasoning options.
193 */
194 configOptions(signal?: AbortSignal): Promise<SessionConfigOption[]> {
195 this.assertActive()
196 return this.modelControl.options(signal)
197 }
198
199 /**
200 * Apply one standard configuration option to later ACP turns.
201 * @param configId - advertised standard option id.
202 * @param value - selected standard option value.
203 * @param signal - optional request cancellation.
204 * @returns the complete resulting option state.
205 */
206 setConfig(configId: string, value: unknown, signal?: AbortSignal): Promise<SessionConfigOption[]> {
207 this.assertActive()
208 return this.modelControl.set(configId, value, signal)
209 }
210
211 /** Resolve topology state off-chain, then serialize its notification without blocking execution updates. */
212 topologyChanged(): void {
213 if (this.closing !== undefined) return
214 void this.modelControl.options()
215 .then((configOptions) => {
216 if (this.closing !== undefined) return
217 const previous = this.outputTail
218 this.outputTail = previous
219 .then(() => this.notify({
220 sessionId: this.agent.session.id,
221 update: { sessionUpdate: 'config_option_update', configOptions },
222 }))
223 /* v8 ignore start -- the bridge notifier contains transport failure. */
224 .catch((error: unknown) => {
225 this.ctx.logger.warn(`acp: config-option update failed: ${errorChain(error)}`)
226 })
227 /* v8 ignore stop */
228 })
229 /* v8 ignore start -- option discovery contains per-provider failure. */
230 .catch((error: unknown) => {
231 this.ctx.logger.warn(`acp: config-option update failed: ${errorChain(error)}`)
232 })
233 /* v8 ignore stop */
234 }
235
236 /**
237 * Admit, enqueue, and settle one prompt at whole-Agent quiescence.
238 * @param params - standard ACP prompt request for this session.
239 * @param imageEnabled - connection capability advertised at initialization.
240 * @param requestSignal - JSON-RPC request cancellation signal.
241 * @returns the correlated standard stop reason after ordered updates drain.
242 */
243 async prompt(
244 params: PromptRequest,
245 imageEnabled: boolean,
246 requestSignal?: AbortSignal,
247 ): Promise<PromptResponse> {
248 this.assertActive()
249 if (this.inflight !== undefined) throw invalidParams('a prompt is already in flight for this session')
250 const completion = Promise.withResolvers<StopReason>()
251 const admission = Promise.withResolvers<void>()
252 const admissionController = new AbortController()
253 const inflight: InflightPrompt = {
254 resolve: completion.resolve,
255 reject: completion.reject,
256 messageId: undefined,
257 messageQueued: false,
258 turn: undefined,
259 endReason: undefined,
260 admissionDone: admission.promise,
261 finishAdmission: admission.resolve,
262 admissionController,
263 cancelRequested: false,
264 settlementStarted: false,
265 outputError: undefined,
266 agentError: undefined,
267 }
268 this.inflight = inflight
269 const onRequestAbort = (): void => { this.cancelPrompt('ACP prompt request cancelled') }
270 requestSignal?.addEventListener('abort', onRequestAbort, { once: true })
271 /* v8 ignore next -- the SDK dispatches a live signal, then notifies abort through its listener. */
272 if (requestSignal?.aborted === true) onRequestAbort()
273 try {
274 let admissionFailure: unknown
275 const promptSelection = this.modelControl.snapshot()
276 try {
277 if (this.ctx.agents.get(this.agent.id) !== this.agent) {
278 throw internalError('prompt was not queued: the agent was disposed outside the bridge')
279 }
280 const content = await admitAcpPrompt(
281 this.ctx,
282 promptSelection,
283 params.prompt,
284 imageEnabled,
285 admissionController.signal,
286 )
287 admissionController.signal.throwIfAborted()
288 if (this.ctx.agents.get(this.agent.id) !== this.agent) {
289 throw internalError('prompt was not queued: the agent was disposed outside the bridge')
290 }
291 const message = createUserMessage({
292 content,
293 source: { kind: 'user' },
294 })
295 inflight.messageId = message.id
296 inflight.messageQueued = true
297 if (promptSelection !== undefined) this.pendingSelections.set(message.id, promptSelection)
298 try {
299 this.agent.followup(message)
300 } catch (error: unknown) {
301 inflight.messageQueued = false
302 this.pendingSelections.delete(message.id)
303 throw error
304 }
305 } catch (error: unknown) {
306 admissionFailure = error
307 } finally {
308 inflight.finishAdmission()
309 }
310
311 if (inflight.cancelRequested) {
312 this.settleAfterQuiescence(inflight)
313 return { stopReason: await completion.promise }
314 }
315 if (admissionFailure !== undefined) {
316 this.inflight = undefined
317 if (admissionFailure instanceof AcpContentError) {
318 throw admissionFailure.kind === 'invalid'
319 ? invalidParams(admissionFailure.message)
320 : internalError(admissionFailure.message)
321 }
322 if (admissionFailure instanceof RequestError) throw admissionFailure
323 throw internalError(`prompt was not queued: ${(admissionFailure as Error).message}`)
324 }
325
326 this.settleAfterQuiescence(inflight)
327 return { stopReason: await completion.promise }
328 } finally {
329 requestSignal?.removeEventListener('abort', onRequestAbort)
330 }
331 }
332
333 /** Cancel the active prompt, or autonomous work when no ACP prompt exists. */
334 cancel(): void {
335 const inflight = this.inflight
336 this.cancelPrompt('ACP prompt cancelled')
337 if (inflight === undefined) this.agent.cancel({ kind: 'user' })
338 }
339
340 /**
341 * Process one durable event and enqueue its standard ACP projections.
342 * @param session - exact event-owning Session.
343 * @param event - committed durable event.
344 */
345 onSessionEvent(session: Session, event: SessionEvent): void {
346 try {
347 if (event.type === 'assistant/message') {
348 const inflight = this.inflight?.turn === event.data.turn ? this.inflight : undefined
349 const previous = this.outputTail
350 const delivery = previous.then(async () => {
351 for (const update of await assistantUpdates(this.ctx, session, event)) {
352 await this.notify({ sessionId: this.agent.session.id, update })
353 }
354 })
355 this.outputTail = delivery.catch((error: unknown) => {
356 const failure = error as Error
357 if (inflight !== undefined) inflight.outputError ??= failure
358 this.ctx.logger.warn(`acp: assistant output conversion failed: ${errorChain(error)}`)
359 })
360 } else if (event.type === 'tool/call') {
361 const previous = this.outputTail
362 this.outputTail = previous
363 .then(() => this.notify({ sessionId: this.agent.session.id, update: toolCallUpdate(event) }))
364 /* v8 ignore start -- the bridge notifier contains transport rejection. */
365 .catch((error: unknown) => {
366 this.ctx.logger.warn(`acp: tool-call update delivery failed: ${errorChain(error)}`)
367 })
368 /* v8 ignore stop */
369 } else if (event.type === 'tool/result') {
370 const previous = this.outputTail
371 this.outputTail = previous
372 .then(async () => this.notify({
373 sessionId: this.agent.session.id,
374 update: await toolResultUpdate(this.ctx, event),
375 }))
376 /* v8 ignore start -- supplemental-content conversion failure is contained and cannot fail Agent work. */
377 .catch((error: unknown) => {
378 this.ctx.logger.warn(`acp: tool-result update delivery failed: ${errorChain(error)}`)
379 })
380 /* v8 ignore stop */
381 }
382 } finally {
383 const inflight = this.inflight
384 if (inflight !== undefined && event.type === 'turn/end' && inflight.turn === event.data.turn) {
385 inflight.endReason = event.data.reason
386 }
387 if (event.type === 'turn/end') this.modelControl.releaseTurn(event.data.turn)
388 }
389 }
390
391 /**
392 * Correlate an accepted user message with its Agent turn and pinned route.
393 * @param message - claimed durable inbox message.
394 * @param turn - allocated Agent turn.
395 */
396 onInboxClaimed(message: UserMessage, turn: number): void {
397 if (this.inflight !== undefined && this.inflight.messageId === message.id) this.inflight.turn = turn
398 const selection = this.pendingSelections.get(message.id)
399 this.pendingSelections.delete(message.id)
400 if (selection !== undefined) this.modelControl.pinTurn(turn, selection)
401 }
402
403 /**
404 * Correlate an Agent interval failure with the active ACP prompt.
405 * @param turn - failed turn number.
406 * @param error - original same-process failure.
407 */
408 onAgentError(turn: number, error: unknown): void {
409 const inflight = this.inflight
410 if (inflight === undefined || !inflight.messageQueued) return
411 // AgentLoop balances an in-turn failure with durable turn/end; settlement
412 // reads that exact error reason. This slot records interval failures outside it.
413 if (inflight.turn === turn) return
414 inflight.agentError = new Error(errorChain(error))
415 this.settleAfterQuiescence(inflight)
416 }
417
418 /** Await every update queued before this call. */
419 drainUpdates(): Promise<void> {
420 return this.outputTail
421 }
422
423 /**
424 * Cancel, drain, flush, and dispose this session once.
425 * @param detail - cancellation detail for any prompt still in admission.
426 * @returns the shared quiescent teardown promise.
427 */
428 close(detail: string): Promise<void> {
429 if (this.closing !== undefined) return this.closing
430 this.closing = (async () => {
431 const failures: unknown[] = []
432 const inflight = this.inflight
433 this.cancelPrompt(detail)
434 if (inflight === undefined || !inflight.messageQueued) this.agent.cancel({ kind: 'user' })
435 try {
436 await inflight?.admissionDone
437 await this.agent.whenIdle()
438 await this.outputTail
439 } catch (error: unknown) {
440 failures.push(new Error('ACP session activity drain failed', { cause: error }))
441 }
442 const subagents = this.ctx.get('subagents') as ContinuableDrain | undefined
443 try {
444 await subagents?.drainContinuableDescendants([this.agent])
445 } catch (error: unknown) {
446 this.ctx.logger.warn(`acp: continuable subagent teardown failed: ${errorChain(error)}`)
447 failures.push(new Error('continuable subagent teardown failed', { cause: error }))
448 }
449 try {
450 await this.ctx.sessions.flush(this.agent.session)
451 } catch (error: unknown) {
452 failures.push(new Error('ACP session persistence flush failed', { cause: error }))
453 }
454 try {
455 await this.disposeAgent()
456 } catch (error: unknown) {
457 failures.push(error)
458 }
459 this.pendingSelections.clear()
460 if (failures.length === 1) throw failures[0]
461 /* v8 ignore start -- independent teardown failures can aggregate only under multiple simultaneous provider faults. */
462 if (failures.length > 1) {
463 throw new AggregateError(failures, `ACP session teardown failed: ${failures.map(errorChain).join('; ')}`)
464 }
465 /* v8 ignore stop */
466 })()
467 return this.closing
468 }
469
470 private assertActive(): void {
471 if (this.closing !== undefined) throw invalidParams(`session is closing: ${this.agent.session.id}`)
472 }
473
474 private cancelPrompt(detail: string): void {
475 const inflight = this.inflight
476 if (inflight === undefined) return
477 inflight.cancelRequested = true
478 inflight.admissionController.abort(new Error(detail))
479 this.settleAfterQuiescence(inflight)
480 if (inflight.messageQueued) this.agent.cancel({ kind: 'user' })
481 }
482
483 private settleAfterQuiescence(inflight: InflightPrompt): void {
484 if (inflight.settlementStarted) return
485 inflight.settlementStarted = true
486 void (async () => {
487 await inflight.admissionDone
488 if (inflight.messageQueued) {
489 await this.agent.whenIdle()
490 await this.outputTail
491 }
492 /* v8 ignore next -- this prompt owns the slot until this exact settlement clears it. */
493 if (this.inflight !== inflight) return
494 this.inflight = undefined
495 if (inflight.cancelRequested) {
496 inflight.resolve('cancelled')
497 return
498 }
499 if (inflight.outputError !== undefined) {
500 inflight.reject(internalError(`assistant output delivery failed: ${inflight.outputError.message}`))
501 return
502 }
503 if (inflight.agentError !== undefined) {
504 inflight.reject(internalError(`turn failed: ${inflight.agentError.message}`))
505 return
506 }
507 const end = inflight.endReason
508 if (end === undefined) {
509 inflight.resolve('cancelled')
510 } else if (end.kind === 'error') {
511 inflight.reject(internalError(`turn failed: ${end.error.message}`))
512 } else {
513 inflight.resolve(turnEndToStopReason(end))
514 }
515 })()
516 /* v8 ignore start -- admissionDone only resolves; idle/output gates contain their own failures. */
517 .catch((error: unknown) => {
518 if (this.inflight !== inflight) return
519 this.inflight = undefined
520 inflight.reject(internalError(`prompt settlement failed: ${errorChain(error)}`))
521 })
522 /* v8 ignore stop */
523 }
524}