index.ts 7.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194
  1. /**
  2. * Shared driver for in-process subagent providers. The agent factory's
  3. * creation transaction owns unpublished setup and rollback; after publication
  4. * the returned AgentHandle is the one quiescent lifecycle owner held by the
  5. * provider's caller.
  6. *
  7. * @module @deepseek-ai/dsh-subagent-inprocess
  8. */
  9. import { randomUUID } from 'node:crypto'
  10. import type { Context } from 'cordis'
  11. import type { Agent, AgentOptions } from '@deepseek-ai/dsh-agent'
  12. import { findLastMessageTurnEnd, SessionId, type SessionEvent, type TurnEndReason } from '@deepseek-ai/dsh-session'
  13. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  14. import { assertSubagentMaxDepth, delegationDepthOf } from '@deepseek-ai/dsh-subagent'
  15. import type { SubagentResult, SubagentRun, SubagentStartRequest, SubagentStopReason } from '@deepseek-ai/dsh-subagent'
  16. import {
  17. attachStructuredRuntime,
  18. type StructuredAttachment,
  19. } from './structured.ts'
  20. export {
  21. STRUCTURED_OUTPUT_TOOL,
  22. STRUCTURED_OUTPUT_INSTRUCTION,
  23. } from './structured.ts'
  24. /** Thrown when starting a child would exceed the requested depth cap. */
  25. class SubagentDepthError extends Error {
  26. constructor(public readonly attemptedDepth: number, public readonly maxDepth: number) {
  27. super(`subagent depth ${attemptedDepth} exceeds maxDepth ${maxDepth}`)
  28. this.name = 'SubagentDepthError'
  29. }
  30. }
  31. /** Map a session turn outcome to the subagent seam's terminal vocabulary. */
  32. function toStopReason(reason: TurnEndReason | undefined): SubagentStopReason {
  33. switch (reason?.kind) {
  34. case 'completed':
  35. return 'completed'
  36. case 'max-tokens':
  37. return 'max-tokens'
  38. case 'aborted':
  39. return 'aborted'
  40. case 'error':
  41. case 'disposed':
  42. case 'interrupted':
  43. default:
  44. return 'error'
  45. }
  46. }
  47. /** Extra inputs the spawn and fork providers supply to the shared driver. */
  48. export interface InProcessRunOptions {
  49. /** Completed-turn seed for fork, or undefined for a fresh spawn. */
  50. readonly seed?: SessionEvent[]
  51. }
  52. /** Error used when cancellation wins before the child publication boundary. */
  53. function prePublicationAbort(): Error {
  54. return new Error('subagent request was aborted before child publication')
  55. }
  56. /**
  57. * Establish and drive one in-process child. Fulfillment means the agent is
  58. * already published in the registry; rejection means the agent factory's
  59. * creation transaction and any partially-created child have reached quiescence.
  60. * @param request - the trusted typed start request, including its required signal.
  61. * @param options - the optional fork seed.
  62. * @returns a ready holder-owned run.
  63. */
  64. export async function startInProcessRun(
  65. request: SubagentStartRequest,
  66. options: InProcessRunOptions,
  67. ): Promise<SubagentRun> {
  68. assertSubagentMaxDepth(request.maxDepth)
  69. if (request.signal.aborted) throw prePublicationAbort()
  70. const parent = request.parent
  71. const childDepth = delegationDepthOf(parent) + 1
  72. if (!Number.isSafeInteger(childDepth)) {
  73. throw new RangeError('subagent child depth exceeds the safe-integer range')
  74. }
  75. if (request.maxDepth !== undefined && childDepth > request.maxDepth) {
  76. throw new SubagentDepthError(childDepth, request.maxDepth)
  77. }
  78. const childId = SessionId(randomUUID())
  79. const seedLength = options.seed?.length ?? 0
  80. const parentHeader = parent.session.header
  81. const parentProvider = parent.options.provider
  82. const parentModel = parent.options.model
  83. const agentOptions: AgentOptions = {
  84. ...parentProvider !== undefined ? { provider: parentProvider } : {},
  85. ...parentModel !== undefined ? { model: parentModel } : {},
  86. ...request.agentOptions,
  87. subagentDepth: childDepth,
  88. }
  89. let structured: StructuredAttachment | undefined
  90. const setup = (childCtx: Context): void => {
  91. if (request.persona !== undefined) {
  92. childCtx.systemPrompt.section({ name: 'deployment:persona', order: 0, text: request.persona })
  93. }
  94. if (request.toolFilter !== undefined) childCtx.tools.restrict(request.toolFilter)
  95. if (request.outputSchema !== undefined) {
  96. structured = attachStructuredRuntime(childCtx, request.outputSchema)
  97. }
  98. }
  99. const flags = { cancelled: false }
  100. const handle = await parent.ctx.agents.create({
  101. sessionId: childId,
  102. meta: {
  103. ...parentHeader.cwd !== undefined ? { cwd: parentHeader.cwd } : {},
  104. parentSession: parentHeader.id,
  105. // Durable: the recursion budget must survive persistence and resume.
  106. delegationDepth: childDepth,
  107. ...seedLength > 0 ? { seedLength } : {},
  108. },
  109. ...options.seed !== undefined ? { seed: options.seed } : {},
  110. agentOptions,
  111. signal: request.signal,
  112. setup,
  113. })
  114. const child = handle.agent
  115. // Agent creation detaches its creation-only abort listener before returning.
  116. // Close the narrow handoff race before installing the live-run listener.
  117. // Static analysis does not model the abort that may land between the
  118. // factory's listener detachment and this continuation.
  119. // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition
  120. if (request.signal.aborted) {
  121. flags.cancelled = true
  122. await handle.dispose()
  123. throw prePublicationAbort()
  124. }
  125. const onAbort = (): void => {
  126. flags.cancelled = true
  127. child.cancel({ kind: 'parent' })
  128. }
  129. request.signal.addEventListener('abort', onAbort, { once: true })
  130. const result: Promise<SubagentResult> = (async () => {
  131. try {
  132. child.followup(request.prompt)
  133. await child.whenIdle()
  134. return readResult(
  135. child,
  136. seedLength,
  137. flags.cancelled,
  138. structured ? { captured: structured.captured() } : undefined,
  139. )
  140. } finally {
  141. request.signal.removeEventListener('abort', onAbort)
  142. }
  143. })()
  144. return {
  145. id: childId,
  146. localAgent: child,
  147. result,
  148. dispose(): Promise<void> {
  149. request.signal.removeEventListener('abort', onAbort)
  150. flags.cancelled = true
  151. return handle.dispose()
  152. },
  153. }
  154. }
  155. /** Read one settled child's result from events after its optional fork seed. */
  156. function readResult(
  157. child: Agent,
  158. seedLength: number,
  159. cancelled: boolean,
  160. structured?: { captured?: { value: unknown } | undefined },
  161. ): SubagentResult {
  162. const own = child.session.events.slice(seedLength)
  163. const lastMessage = own.findLast((event): event is SessionEvent<'assistant/message'> => event.type === 'assistant/message')
  164. const lastEnd = findLastMessageTurnEnd(own)
  165. const output: ContentBlock[] = lastMessage?.data.content ?? []
  166. const recorded = toStopReason(lastEnd?.data.reason)
  167. // Disposal can tear the owner down before the loop records its ordinary
  168. // `aborted` end, yielding `disposed` instead. A requested cancellation owns
  169. // every non-completed in-flight outcome; a turn already completed stays so.
  170. const stopReason: SubagentStopReason = cancelled && recorded !== 'completed'
  171. ? 'aborted'
  172. : recorded
  173. if (structured !== undefined) {
  174. if (structured.captured !== undefined) {
  175. return { output, structured: structured.captured.value, stopReason }
  176. }
  177. if (stopReason === 'completed') return { output, stopReason: cancelled ? 'aborted' : 'error' }
  178. }
  179. return { output, stopReason }
  180. }