index.ts 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235
  1. /**
  2. * Shared driver for in-process ONE-SHOT 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. * Continuable children never come through here: the continuation manager
  8. * composes and drives them directly, so this driver owns exactly one turn with
  9. * one result.
  10. *
  11. * @module @deepseek-ai/dsh-subagent-inprocess
  12. */
  13. import { randomUUID } from 'node:crypto'
  14. import type { Context } from 'cordis'
  15. import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent'
  16. import { findLastMessageTurnEnd, SessionId, type SessionEvent, type TurnEndReason } from '@deepseek-ai/dsh-session'
  17. import { createUserMessage, type ContentBlock } from '@deepseek-ai/dsh-llm'
  18. import {
  19. applyChildComposition,
  20. assertSubagentMaxDepth,
  21. childSessionMeta,
  22. resolveChildAgentOptions,
  23. resolveChildDepth,
  24. } from '@deepseek-ai/dsh-subagent'
  25. import type {
  26. ResolvedSubagentStartRequest,
  27. SubagentDescriptorData,
  28. SubagentResult,
  29. SubagentRun,
  30. SubagentStopReason,
  31. } from '@deepseek-ai/dsh-subagent'
  32. // Type-only: make `ctx.get('sandboxPolicy')` / `ctx.get('approval')` resolve
  33. // to the policy services when composed — the driver consumes both
  34. // opportunistically (the documented `ctx.get` pattern), never as a hard dep.
  35. import type {} from '@deepseek-ai/dsh-sandbox-policy'
  36. import type {} from '@deepseek-ai/dsh-user-approval'
  37. import {
  38. attachStructuredRuntime,
  39. type StructuredAttachment,
  40. } from './structured.ts'
  41. export {
  42. STRUCTURED_OUTPUT_TOOL,
  43. STRUCTURED_OUTPUT_INSTRUCTION,
  44. } from './structured.ts'
  45. /** Map a session turn outcome to the subagent seam's terminal vocabulary. */
  46. function toStopReason(reason: TurnEndReason | undefined): SubagentStopReason {
  47. switch (reason?.kind) {
  48. case 'completed':
  49. return 'completed'
  50. case 'max-tokens':
  51. return 'max-tokens'
  52. case 'aborted':
  53. return 'aborted'
  54. case 'error':
  55. case 'interrupted':
  56. default:
  57. return 'error'
  58. }
  59. }
  60. /** Extra inputs the spawn and fork providers supply to the shared driver. */
  61. export interface InProcessRunOptions {
  62. /** Completed-turn seed for fork, or undefined for a fresh spawn. */
  63. readonly seed?: SessionEvent[]
  64. }
  65. /** Error used when cancellation wins before the child publication boundary. */
  66. function prePublicationAbort(): Error {
  67. return new Error('subagent request was aborted before child publication')
  68. }
  69. /** Append one one-shot descriptor inside the child's initial turn before its first request. */
  70. function attachDescriptorAppend(childCtx: Context, descriptor: SubagentDescriptorData): void {
  71. let appended = false
  72. childCtx.on('agent/pre-step', async ({ agent }, next) => {
  73. const decision = await next()
  74. if (!appended && decision.kind === 'enter') {
  75. appended = true
  76. agent.session.append('subagent/descriptor', descriptor)
  77. }
  78. return decision
  79. })
  80. }
  81. /**
  82. * Establish and drive one in-process one-shot child. Fulfillment means the agent
  83. * is already published in the registry and transfers its turn, cancellation,
  84. * and disposal work through the returned run. Rejection means the agent
  85. * factory's unpublished creation transaction reached quiescence without
  86. * publishing a child. Every start appends its resolved descriptor inside the
  87. * child's initial turn.
  88. * @param request - the trusted typed start request, including its required signal.
  89. * @param options - the optional fork seed.
  90. * @returns a published holder-owned run.
  91. */
  92. export async function startInProcessRun(
  93. request: ResolvedSubagentStartRequest,
  94. options: InProcessRunOptions,
  95. ): Promise<SubagentRun> {
  96. assertSubagentMaxDepth(request.maxDepth)
  97. if (request.signal.aborted) throw prePublicationAbort()
  98. const parent = request.parent
  99. const childDepth = resolveChildDepth(parent, request.maxDepth)
  100. const childId = SessionId(randomUUID())
  101. const seed = options.seed
  102. const activationBoundary = seed?.length ?? 0
  103. // Capture before the first await: a later parent switch belongs to the
  104. // parent's future.
  105. const inheritedMode = parent.ctx.get('sandboxPolicy')?.overrideOf(parent.session)
  106. const inheritedPolicy = parent.ctx.get('approval')?.overrideOf(parent.session)
  107. let structured: StructuredAttachment | undefined
  108. const setup = (childCtx: Context): void => {
  109. // Inherited overrides land on the child's own log, so its effective policy
  110. // is reconstructable from that log alone.
  111. const childSession = (childCtx.agent as Agent).session
  112. if (inheritedMode !== undefined) {
  113. childSession.append('sandbox/mode', { mode: inheritedMode, source: 'delegation' })
  114. }
  115. if (inheritedPolicy !== undefined) {
  116. childSession.append('approval/policy', { policy: inheritedPolicy, source: 'delegation' })
  117. }
  118. applyChildComposition(childCtx, {
  119. persona: request.persona,
  120. toolFilter: request.toolFilter,
  121. })
  122. if (request.outputSchema !== undefined) {
  123. structured = attachStructuredRuntime(childCtx, request.outputSchema)
  124. }
  125. attachDescriptorAppend(childCtx, request.descriptor)
  126. }
  127. const handle = await parent.ctx.agents.create({
  128. sessionId: childId,
  129. meta: childSessionMeta(parent, childDepth, activationBoundary),
  130. ...seed !== undefined ? { seed } : {},
  131. agentOptions: resolveChildAgentOptions(parent, request.agentOptions, childDepth),
  132. signal: request.signal,
  133. setup,
  134. })
  135. return drivePublishedRun(
  136. handle,
  137. request.signal,
  138. request.prompt,
  139. childId,
  140. activationBoundary,
  141. structured,
  142. )
  143. }
  144. /**
  145. * Wrap a published child in the single run lifecycle that owns signal handoff,
  146. * one turn, result settlement, and quiescent disposal.
  147. */
  148. function drivePublishedRun(
  149. handle: AgentHandle,
  150. signal: AbortSignal,
  151. prompt: ContentBlock[],
  152. childId: SessionId,
  153. boundary: number,
  154. structured: StructuredAttachment | undefined,
  155. ): SubagentRun {
  156. const child = handle.agent
  157. const flags = { cancelled: false }
  158. const onAbort = (): void => {
  159. flags.cancelled = true
  160. child.cancel({ kind: 'parent' })
  161. }
  162. signal.addEventListener('abort', onAbort, { once: true })
  163. // Agent creation detaches its creation-only listener before returning. The
  164. // post-registration check closes that handoff without treating an already
  165. // published child as a failed start.
  166. if (signal.aborted) onAbort()
  167. const result: Promise<SubagentResult> = (async () => {
  168. try {
  169. if (!flags.cancelled) {
  170. child.followup(createUserMessage({ content: prompt, source: { kind: 'user' } }))
  171. await child.whenIdle()
  172. }
  173. return readResult(
  174. child,
  175. boundary,
  176. flags.cancelled,
  177. structured ? { captured: structured.captured() } : undefined,
  178. )
  179. } finally {
  180. signal.removeEventListener('abort', onAbort)
  181. }
  182. })()
  183. return {
  184. id: childId,
  185. localAgent: child,
  186. result,
  187. async dispose(): Promise<void> {
  188. signal.removeEventListener('abort', onAbort)
  189. flags.cancelled = true
  190. const settlements = await Promise.allSettled([handle.dispose(), result])
  191. const disposal = settlements[0]
  192. // The result channel owns run faults; disposal reports only failure to
  193. // release the published handle after both operations settle.
  194. if (disposal.status === 'rejected') throw disposal.reason
  195. },
  196. }
  197. }
  198. /** Read one settled child's result from events after its activation boundary. */
  199. function readResult(
  200. child: Agent,
  201. boundary: number,
  202. cancelled: boolean,
  203. structured?: { captured?: { value: unknown } | undefined },
  204. ): SubagentResult {
  205. const own = child.session.events.slice(boundary)
  206. const lastMessage = own.findLast((event): event is SessionEvent<'assistant/message'> => event.type === 'assistant/message')
  207. const lastEnd = findLastMessageTurnEnd(own)
  208. const output: ContentBlock[] = lastMessage?.data.message.content ?? []
  209. const recorded = toStopReason(lastEnd?.data.reason)
  210. // Disposal can tear the owner down before the loop records its ordinary
  211. // `aborted` end, yielding `disposed` instead.
  212. const stopReason: SubagentStopReason = cancelled && recorded !== 'completed' ? 'aborted' : recorded
  213. if (structured !== undefined) {
  214. if (structured.captured !== undefined) {
  215. return { output, structured: structured.captured.value, stopReason }
  216. }
  217. if (stopReason === 'completed') return { output, stopReason: cancelled ? 'aborted' : 'error' }
  218. }
  219. return { output, stopReason }
  220. }