返回源码地图

packages/subagent/subagent/src/continuation.ts

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

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

1/**
2 * Continuable-subagent orchestration behind `ctx.subagents`: stable child ids,
3 * descriptor persistence, provider preparation, cold resume, authorization,
4 * and message routing. {@link ContinuableActivationRegistry} owns the mutable
5 * process-local Activation graph and its settlement and disposal lifecycle.
6 *
7 * A continuable child has one durable Session and at most one process-local
8 * Activation. The Agent inbox is the only turn queue, so this manager owns
9 * durable orchestration while the Agent loop owns all turn ordering and
10 * execution. No continuable path creates a Task or an intermediate
11 * result-bearing wrapper.
12 *
13 * @module @deepseek-ai/dsh-subagent
14 */
15
16import { randomUUID } from 'node:crypto'
17import type { Context } from '@deepseek-ai/cordis'
18import type { Agent } from '@deepseek-ai/dsh-agent'
19import { brandString } from '@deepseek-ai/dsh-brand'
20import { ReasoningEffortId, contentHasImage, createUserMessage } from '@deepseek-ai/dsh-llm'
21import type { ContentBlock, MessageId, MessageSource } from '@deepseek-ai/dsh-llm'
22import { SessionLogOffset } from '@deepseek-ai/dsh-session'
23import type { SessionId } from '@deepseek-ai/dsh-session'
24import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
25import type { SessionObservation, SessionQueryEngine } from '@deepseek-ai/dsh-session-query'
26import {
27 childSessionMeta,
28 captureDelegatedPolicyOverrides,
29 resolveChildAgentOptions,
30 resolveChildDepth,
31} from './child-agent.ts'
32import {
33 ContinuableActivationRegistry,
34} from './continuation-activation.ts'
35import type { Activation } from './continuation-activation.ts'
36import {
37 createAgentMessage,
38 withContinuableReturnGuidance,
39} from './continuation-messages.ts'
40import { assertSubagentMaxDepth } from './depth.ts'
41import { foldSubagentDescriptor, snapshotSubagentDescriptor } from './descriptor.ts'
42import { establishCatalogChild } from './catalog.ts'
43import { SubagentError } from './error.ts'
44import { isAdjacentAgentSendMessageTool } from './internal.ts'
45import type { ActivationObserver } from './lifecycle.ts'
46import type {
47 ContinuableCreateRequest,
48 ContinuableCreateSpec,
49 ContinuableStart,
50 ContinuableStartSpec,
51 SubagentInterruptAuthority,
52 SubagentSendMessageOptions,
53} from './types.ts'
54
55/** Inputs shared by model steering and human prompt delivery. */
56type ChildDeliveryOptions =
57 | {
58 readonly delivery: 'steer'
59 /**
60 * A provided host source is preserved on the user message; omission attributes
61 * an adjacent-Agent message to the parent.
62 */
63 readonly source?: MessageSource
64 readonly signal: AbortSignal
65 }
66 | { readonly delivery: 'queue'; readonly source: MessageSource; readonly signal: AbortSignal }
67
68/** Package-private hooks supplied by the owning service. */
69interface ContinuationHost {
70 /** Resolve one provider's detached continuable-creation contribution. */
71 prepareContinuable(name: string, request: ContinuableCreateRequest): Promise<ContinuableCreateSpec>
72 /** Build the lifecycle observer for one Activation residency epoch. */
73 observeActivation(provider: string, childId: SessionId, parent: Agent): ActivationObserver
74}
75
76/**
77 * The continuable-subagent orchestration service behind `ctx.subagents`. Tool
78 * schema and host adapters are consumers of this one contract; foreground
79 * one-shot delegation keeps calling `ctx.subagents.start()` and never enters
80 * this lifecycle.
81 */
82export class SubagentContinuationManager {
83 private readonly activations: ContinuableActivationRegistry
84
85 constructor(
86 private readonly ctx: Context,
87 private readonly host: ContinuationHost,
88 maxActiveSubagents: () => number,
89 ) {
90 this.activations = new ContinuableActivationRegistry(
91 ctx,
92 (provider, childId, parent) => host.observeActivation(provider, childId, parent),
93 maxActiveSubagents,
94 )
95 }
96
97 /**
98 * Start one continuable background child and resolve at initial inbox acceptance.
99 * Every earlier failure disposes any created handle and rolls back Activation
100 * and parent ownership without returning either id.
101 * @param spec - provider, delegation request, and caller cancellation.
102 * @returns the durable child id and accepted initial prompt message id.
103 */
104 async startContinuable(spec: ContinuableStartSpec): Promise<ContinuableStart> {
105 const request = spec.request
106 const parent = request.parent
107 this.activations.assertAdmitting(parent)
108 const persistence = this.requirePersistence()
109 assertSubagentMaxDepth(request.maxDepth)
110 const childId = spec.childId ?? brandString<SessionId>(randomUUID())
111 this.activations.assertChildIdAvailable(childId)
112 const childDepth = resolveChildDepth(parent, request.maxDepth)
113 // Snapshot before any await: invalid descriptor JSON rejects the call
114 // before a child exists, and the detached value is what reaches the log.
115 const agentOptions = resolveChildAgentOptions(parent, request.agentOptions, childDepth)
116 const agentProvider = agentOptions.provider
117 const agentModel = agentOptions.model
118 const agentReasoningEffort = agentOptions.reasoningEffort
119 const descriptor = snapshotSubagentDescriptor({
120 mode: 'continuable',
121 provider: spec.provider,
122 label: spec.label,
123 ...agentProvider !== undefined ? { agentProvider } : {},
124 ...agentModel !== undefined ? { agentModel } : {},
125 ...agentReasoningEffort !== undefined ? { agentReasoningEffort } : {},
126 ...request.persona !== undefined ? { persona: request.persona } : {},
127 ...request.toolFilter !== undefined ? { toolFilter: request.toolFilter } : {},
128 })
129 // Capture before the first await: a later parent switch belongs to the
130 // parent's future, not to this child.
131 const delegatedPolicies = captureDelegatedPolicyOverrides(parent)
132
133 // An idle continuation-managed parent must not settle while a caller is
134 // still creating its child. A turn-scoped delegation does not need this,
135 // but the service is also callable outside a turn.
136 const releaseHold = this.activations.holdOwnership(parent, childId)
137 try {
138 const prepared = await this.host.prepareContinuable(spec.provider, {
139 sessionId: childId,
140 parent,
141 signal: spec.signal,
142 })
143 spec.signal.throwIfAborted()
144 this.activations.assertAdmitting(parent)
145
146 const inheritedEventCount = SessionLogOffset(prepared.seed?.length ?? 0)
147 const seed = prepared.seed
148 const messageId = await this.activations.locks.run(childId, async () => {
149 spec.signal.throwIfAborted()
150 this.activations.assertAdmitting(parent)
151 this.activations.assertChildIdAvailable(childId)
152 if (spec.childId !== undefined) {
153 const persisted = await persistence.stat(childId, { signal: spec.signal })
154 spec.signal.throwIfAborted()
155 this.activations.assertAdmitting(parent)
156 this.activations.assertChildIdAvailable(childId)
157 if (persisted !== undefined) {
158 throw new SubagentError(`subagent "${childId}" already exists`, 'DUPLICATE_CHILD')
159 }
160 }
161 const activation = await this.activations.materialize({
162 childId,
163 provider: spec.provider,
164 parent,
165 create: {
166 seed,
167 meta: childSessionMeta(parent, childDepth, prepared.seed !== undefined),
168 inheritedEventCount,
169 delegatedPolicies,
170 descriptor,
171 },
172 agentOptions,
173 composition: { persona: request.persona, toolFilter: request.toolFilter },
174 signal: spec.signal,
175 })
176 const childHeader = activation.handle.agent.session.header
177 return await this.submitMaterialized(
178 activation,
179 isAdjacentAgentSendMessageTool(this.ctx.get('tools')?.get('send_message', activation.handle.agent))
180 ? withContinuableReturnGuidance(parent.id, request.prompt)
181 : request.prompt,
182 { source: { kind: 'user' }, signal: spec.signal, delivery: 'queue' },
183 parent,
184 () => { establishCatalogChild(parent.session, childHeader, descriptor) },
185 )
186 })
187 return { childId, messageId }
188 } catch (error: unknown) {
189 releaseHold()
190 throw error
191 }
192 }
193
194 /**
195 * Deliver one model-authored message to a direct continuable child or to the
196 * sender's direct parent. A missing direct child cold-resumes through the
197 * ordinary continuation lifecycle.
198 * @param sender - exact live Agent authorizing and originating the message.
199 * @param targetId - durable direct-parent or direct-child session id.
200 * @param content - model-authored content to deliver.
201 * @param options - caller cancellation before acceptance.
202 * @returns the accepted message's inbox id.
203 */
204 async sendMessage(
205 sender: Agent,
206 targetId: SessionId,
207 content: ContentBlock[],
208 options: SubagentSendMessageOptions,
209 ): Promise<MessageId> {
210 if (this.ctx.agents.get(sender.id) !== sender) {
211 throw new SubagentError(
212 'message delivery requires the exact live sender agent',
213 'UNAUTHORIZED',
214 )
215 }
216 this.activations.assertAdmitting(sender)
217 const senderActivation = this.activations.get(sender.id)
218 if (senderActivation !== undefined
219 && senderActivation.handle.agent === sender
220 && senderActivation.parentSession === targetId) {
221 options.signal.throwIfAborted()
222 return this.sendToParent(senderActivation, sender, content)
223 }
224 if (sender.session.header.parentSession === targetId) {
225 throw new SubagentError(
226 `agent "${sender.id}" is not a resident continuable child and cannot send to parent "${targetId}"`,
227 'UNAUTHORIZED',
228 )
229 }
230 return this.deliverToChild(sender, targetId, content, {
231 signal: options.signal,
232 delivery: 'steer',
233 })
234 }
235
236 /**
237 * Queue one human-authored prompt as a distinct direct-child turn.
238 * @param parent - exact live direct parent authorizing delivery.
239 * @param childId - durable direct-child session id.
240 * @param content - model-visible prompt blocks.
241 * @param source - durable attribution for the human prompt.
242 * @param signal - caller cancellation before inbox acceptance.
243 * @returns the accepted durable message id.
244 */
245 async queuePrompt(
246 parent: Agent,
247 childId: SessionId,
248 content: ContentBlock[],
249 source: MessageSource,
250 signal: AbortSignal,
251 ): Promise<MessageId> {
252 return this.deliverToChild(parent, childId, content, { source, signal, delivery: 'queue' })
253 }
254
255 /**
256 * Steer one host-authored prompt to a direct continuable child.
257 * @param parent - exact live direct parent authorizing delivery.
258 * @param childId - durable direct-child session id.
259 * @param content - model-visible prompt blocks.
260 * @param source - durable attribution for the host prompt.
261 * @param signal - caller cancellation before inbox acceptance.
262 * @returns the accepted durable message id.
263 */
264 async steerPrompt(
265 parent: Agent,
266 childId: SessionId,
267 content: ContentBlock[],
268 source: MessageSource,
269 signal: AbortSignal,
270 ): Promise<MessageId> {
271 return this.deliverToChild(parent, childId, content, { source, signal, delivery: 'steer' })
272 }
273
274 /** Route one parent-originated delivery through residency and cold resume. */
275 private async deliverToChild(
276 parent: Agent,
277 childId: SessionId,
278 content: ContentBlock[],
279 options: ChildDeliveryOptions,
280 ): Promise<MessageId> {
281 this.activations.assertAdmitting(parent)
282 const releaseHold = this.activations.holdOwnership(parent, childId)
283 try {
284 return await this.deliverFollowup(parent, childId, content, options)
285 } catch (error: unknown) {
286 releaseHold()
287 throw error
288 }
289 }
290
291 /** The delivery loop behind {@link deliverToChild}, run under the parent hold. */
292 private async deliverFollowup(
293 parent: Agent,
294 childId: SessionId,
295 content: ContentBlock[],
296 options: ChildDeliveryOptions,
297 ): Promise<MessageId> {
298 while (true) {
299 const live = await this.activations.locks.run(childId, async () => {
300 const activation = this.activations.get(childId)
301 if (activation === undefined) return this.coldResume(parent, childId, content, options)
302 const disposal = activation.inbox.closing
303 /* v8 ignore next 3 -- the send-versus-dispose cutoff needs a delivery to
304 * observe the transaction inside the same critical section that opened it. */
305 if (disposal !== undefined) {
306 return disposal.then(() => undefined, () => undefined)
307 }
308 if (contentHasImage(content)) {
309 await this.assertImageCapable(activation.handle.agent, options.signal)
310 if (activation.inbox.closing !== undefined) {
311 await Promise.allSettled([activation.inbox.closing])
312 return undefined
313 }
314 }
315 const messageId = this.submitAdmitted(activation, content, options, parent)
316 activation.announced = true
317 return messageId
318 })
319 /* v8 ignore start -- only a delivery that lost the disposal cutoff retries. */
320 if (live !== undefined) return live
321 this.activations.assertAdmitting(parent)
322 options.signal.throwIfAborted()
323 /* v8 ignore stop */
324 }
325 }
326
327 /**
328 * Interrupt one live continuable child's current turn. Admission is
329 * synchronous and the cancellation effect is asynchronous. An absent or
330 * already-closing target is an accepted no-op after authority checks.
331 * @param targetSessionId - the durable child session id to interrupt.
332 * @param authority - the human parent address or exact live ancestor Agent.
333 */
334 interrupt(targetSessionId: SessionId, authority: SubagentInterruptAuthority): void {
335 this.activations.interrupt(targetSessionId, authority)
336 }
337
338 /** Deliver one resident continuable child's message to its live direct parent. */
339 private sendToParent(
340 activation: Activation,
341 sender: Agent,
342 content: ContentBlock[],
343 ): MessageId {
344 /* v8 ignore next 6 -- only synchronous re-entrant teardown can open this
345 * transaction between exact-agent authorization and this no-await span. */
346 if (activation.inbox.closing !== undefined) {
347 throw new SubagentError(
348 `subagent "${sender.id}" activation is being disposed; the message was not delivered`,
349 'ACTIVATION_CLOSING',
350 )
351 }
352 const parent = this.ctx.agents.get(activation.parentSession)
353 if (parent === undefined) {
354 throw new SubagentError(
355 'direct parent is not live; the message was not delivered',
356 'PARENT_UNAVAILABLE',
357 )
358 }
359 const message = createAgentMessage(sender, content)
360 this.sendAgentMessage(parent, message)
361 return message.id
362 }
363
364 /** Send one Agent message while translating only the target's own rejection. */
365 private sendAgentMessage(
366 parent: Agent,
367 message: ReturnType<typeof createUserMessage>,
368 ): void {
369 try {
370 this.activations.sendWaking(parent, message, 'steer')
371 } catch (error: unknown) {
372 throw new SubagentError(
373 'direct parent is not live; the message was not delivered',
374 'PARENT_UNAVAILABLE',
375 { cause: error },
376 )
377 }
378 }
379
380 /** Close manager-wide admission and release every live Activation. */
381 async drain(): Promise<void> {
382 await this.activations.drain()
383 }
384
385 /**
386 * Stop only the continuable descendants of exact live host-owned parents.
387 * @param parents - exact live roots whose continuable descendants must stop.
388 */
389 async drainDescendants(parents: readonly Agent[]): Promise<void> {
390 await this.activations.drainDescendants(parents)
391 }
392
393 /**
394 * Release selected resident direct children of one exact live parent.
395 * @param parent - exact live direct parent authorizing the selected release.
396 * @param childIds - durable direct-child ids to release when resident.
397 */
398 async drainChildren(parent: Agent, childIds: readonly SessionId[]): Promise<void> {
399 await this.activations.drainChildren(parent, childIds)
400 }
401
402 /**
403 * Cold-resume a persisted child and submit the waiting turn. The descriptor
404 * supplies every reconstruction input; no subagent provider is dispatched.
405 */
406 private async coldResume(
407 parent: Agent,
408 childId: SessionId,
409 content: ContentBlock[],
410 options: ChildDeliveryOptions,
411 ): Promise<MessageId> {
412 const query = this.requireSessionQuery()
413 let observation: SessionObservation
414 try {
415 observation = await query.observeSession(childId, {
416 signal: options.signal,
417 })
418 } catch (error: unknown) {
419 options.signal.throwIfAborted()
420 throw new SubagentError(`subagent "${childId}" is unavailable`, 'NOT_RESUMABLE', { cause: error })
421 }
422 using source = observation
423 this.activations.assertAdmitting(parent)
424 this.activations.authorizeLineage(parent, childId, source.header.parentSession)
425 const descriptor = foldSubagentDescriptor(
426 source.events.slice(source.inheritedEventCount),
427 )
428 if (descriptor === undefined || descriptor.mode !== 'continuable') {
429 throw new SubagentError(
430 `subagent "${childId}" has no supported continuation state and cannot be resumed; choose a different target`,
431 'NOT_RESUMABLE',
432 )
433 }
434 let activation: Activation
435 try {
436 activation = await this.activations.materialize({
437 childId,
438 provider: descriptor.provider,
439 parent,
440 agentOptions: {
441 ...descriptor.agentProvider !== undefined ? { provider: descriptor.agentProvider } : {},
442 ...descriptor.agentModel !== undefined ? { model: descriptor.agentModel } : {},
443 ...descriptor.agentReasoningEffort !== undefined
444 ? { reasoningEffort: ReasoningEffortId(descriptor.agentReasoningEffort) }
445 : {},
446 },
447 composition: { persona: descriptor.persona, toolFilter: descriptor.toolFilter },
448 signal: options.signal,
449 })
450 } catch (error: unknown) {
451 options.signal.throwIfAborted()
452 if (error instanceof SubagentError) throw error
453 throw new SubagentError(`subagent "${childId}" is unavailable`, 'NOT_RESUMABLE', { cause: error })
454 }
455 return await this.submitMaterialized(activation, content, options, parent)
456 }
457
458 /** Admit a materialized child, commit its creation fact, and release it on failure. */
459 private async submitMaterialized(
460 activation: Activation,
461 content: ContentBlock[],
462 options: ChildDeliveryOptions,
463 parent: Agent,
464 commit?: () => void,
465 ): Promise<MessageId> {
466 try {
467 if (contentHasImage(content)) {
468 await this.assertImageCapable(activation.handle.agent, options.signal)
469 if (activation.inbox.closing !== undefined) {
470 throw new SubagentError(`subagent "${activation.childId}" is closing`, 'ACTIVATION_CLOSING')
471 }
472 }
473 const messageId = this.submitAdmitted(activation, content, options, parent)
474 commit?.()
475 activation.announced = true
476 return messageId
477 } catch (error: unknown) {
478 try {
479 await this.activations.dispose(activation)
480 } catch (cleanupError: unknown) {
481 this.ctx.logger.warn(
482 `subagent continuation: disposal after admission or catalog append failure also failed: ${String(cleanupError)}`,
483 )
484 }
485 throw error
486 }
487 }
488
489 /** Build and submit one message across the final synchronous admission cutoff. */
490 private submitAdmitted(
491 activation: Activation,
492 content: ContentBlock[],
493 options: ChildDeliveryOptions,
494 parent: Agent,
495 ): MessageId {
496 const message = options.source === undefined
497 ? createAgentMessage(parent, content)
498 : createUserMessage({ content, source: options.source })
499 return this.activations.submitAdmitted(
500 activation,
501 message,
502 options.delivery,
503 parent,
504 options.signal,
505 )
506 }
507
508 /** Refuse image content for a child whose fixed model accepts text only. */
509 private async assertImageCapable(
510 agent: Agent,
511 signal: AbortSignal,
512 ): Promise<void> {
513 const { provider, model } = agent.options
514 if (provider === undefined || model === undefined) return
515 const llm = this.ctx.get('llm')
516 /* v8 ignore next -- without an LLM registry, delivery defers to projection. */
517 if (llm === undefined) return
518 const info = await llm.resolveModelInfo(provider, model, signal)
519 if (info.inputModalities !== undefined && !info.inputModalities.includes('image')) {
520 throw new SubagentError(
521 `Model "${model}" does not support image input.`,
522 'MODEL_DOES_NOT_SUPPORT_IMAGES',
523 )
524 }
525 }
526
527 /** Resolve the persistence service continuable children require, or fail loud. */
528 private requirePersistence(): SessionPersistence {
529 const persistence = this.ctx.get('sessionPersistence')
530 if (persistence === undefined) {
531 throw new SubagentError(
532 'continuable subagents require session persistence (load a dsh-session-persistence backend)',
533 'PERSISTENCE_UNAVAILABLE',
534 )
535 }
536 return persistence
537 }
538
539 /** Resolve the Session query service used for cold child observations. */
540 private requireSessionQuery(): SessionQueryEngine {
541 const query = this.ctx.get('sessionQuery')
542 if (query === undefined) {
543 throw new SubagentError(
544 'continuable subagents require session query (load @deepseek-ai/dsh-session-query)',
545 'CONTINUATION_UNAVAILABLE',
546 )
547 }
548 return query
549 }
550}
551
552export default SubagentContinuationManager