run.ts 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  1. /**
  2. * One-shot Claude Code lifecycle: invoke the official Agent SDK, place its
  3. * real CLI process under the shared subprocess owner, map only strict SDK
  4. * success to completion, and dispose to whole-tree quiescence.
  5. *
  6. * @module @deepseek-ai/dsh-subagent-claude-code/run
  7. */
  8. import { randomUUID } from 'node:crypto'
  9. import {
  10. query as officialQuery,
  11. type Options,
  12. type Query,
  13. type SDKMessage,
  14. type SDKResultMessage,
  15. type SpawnOptions,
  16. } from '@anthropic-ai/claude-agent-sdk'
  17. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  18. import { SessionId } from '@deepseek-ai/dsh-session'
  19. import {
  20. settleRunResult,
  21. subprocessRunHandle,
  22. type SubagentResult,
  23. type SubagentRun,
  24. type SubagentStartRequest,
  25. type SubagentStopReason,
  26. } from '@deepseek-ai/dsh-subagent'
  27. import {
  28. scrubbedParentEnv,
  29. type SubprocessHandle,
  30. type SubprocessSpawnSpec,
  31. } from '@deepseek-ai/dsh-subprocess'
  32. import {
  33. claudeSpawnSpec,
  34. ManagedClaudeCodeProcess,
  35. } from './process.ts'
  36. /** Default POSIX grace between subprocess termination tiers. */
  37. export const DEFAULT_DISPOSE_GRACE_MS = 3_000
  38. /* jscpd:ignore-start -- sibling providers intentionally keep product-private
  39. * run inputs and error normalization instead of adding a shared lifecycle owner. */
  40. /** Fully resolved inputs for one official Claude Agent SDK query. */
  41. export interface ClaudeCodeRunSpec {
  42. /** Parent Session workspace supplied to the SDK and real CLI. */
  43. readonly cwd: string
  44. /** Exact native Claude Code executable resolved from the host PATH. */
  45. readonly executable: string
  46. /** Explicit deployment/test environment layered after shared scrubbing. */
  47. readonly env: Record<string, string>
  48. /** Subprocess termination grace passed to the shared process-tree owner. */
  49. readonly disposeGraceMs: number
  50. /** Shared subprocess service spawn operation. */
  51. readonly spawn: (spec: SubprocessSpawnSpec) => SubprocessHandle
  52. /** Diagnostic sink for a post-publication error flattened into a result. */
  53. readonly onError?: (error: Error, stopReason: SubagentStopReason) => void
  54. }
  55. function thrown(value: unknown): Error {
  56. /* v8 ignore next -- typed SDK and subprocess failures reject with Error. */
  57. return value instanceof Error ? value : new Error(String(value))
  58. }
  59. /* jscpd:ignore-end */
  60. /**
  61. * Validate and preserve the one-shot task before crossing the SDK boundary.
  62. * @param prompt - task content accepted from the shared subagent service.
  63. * @returns the exact text sequence as one SDK prompt.
  64. */
  65. export function textTask(prompt: readonly ContentBlock[]): string {
  66. if (prompt.length === 0) {
  67. throw new Error('subagent-claude-code: the one-shot task must contain only text blocks')
  68. }
  69. const texts: string[] = []
  70. for (const block of prompt) {
  71. if (block.type !== 'text') {
  72. throw new Error('subagent-claude-code: the one-shot task must contain only text blocks')
  73. }
  74. texts.push(block.text)
  75. }
  76. if (texts.every(text => text.trim().length === 0)) {
  77. throw new Error('subagent-claude-code: the one-shot task must not be empty')
  78. }
  79. return texts.join('')
  80. }
  81. /**
  82. * Strictly derive the only SDK result that can complete a shared run.
  83. * @param message - an official discriminated result union.
  84. * @returns exact final text for a successful, non-error result.
  85. */
  86. export function successfulResult(message: SDKResultMessage): string {
  87. if (
  88. message.subtype !== 'success'
  89. || message.is_error
  90. || message.result.trim().length === 0
  91. ) {
  92. const detail = message.subtype === 'success'
  93. ? 'success result was marked as an error or contained no answer'
  94. : message.errors.join('; ') || message.subtype
  95. throw new Error(`subagent-claude-code: Claude Code failed: ${detail}`)
  96. }
  97. return message.result
  98. }
  99. /**
  100. * Consume the complete SDK stream and require one strict success plus normal
  101. * iterator completion.
  102. * @param query - published official SDK query.
  103. * @returns the completed shared result.
  104. */
  105. export async function consumeClaudeQuery(
  106. query: AsyncIterable<SDKMessage>,
  107. ): Promise<SubagentResult> {
  108. let answer: string | undefined
  109. for await (const message of query) {
  110. if (message.type !== 'result') continue
  111. answer = successfulResult(message)
  112. }
  113. if (answer === undefined) {
  114. throw new Error('subagent-claude-code: Claude Code ended without a result')
  115. }
  116. return {
  117. output: [{ type: 'text', text: answer }],
  118. stopReason: 'completed',
  119. }
  120. }
  121. /**
  122. * Close the official query, terminate the managed process tree, and wait for
  123. * the subprocess owner to prove it is gone.
  124. * @param query - official SDK query, when creation reached that point.
  125. * @param child - shared-service handle that owns the CLI process tree.
  126. */
  127. export async function disposeClaudeCodeChild(
  128. query: Pick<Query, 'close'> | undefined,
  129. child: SubprocessHandle,
  130. ): Promise<void> {
  131. const failures: Error[] = []
  132. try {
  133. query?.close()
  134. } catch (error: unknown) {
  135. failures.push(thrown(error))
  136. }
  137. if (child.pid > 0) {
  138. child.terminate()
  139. try {
  140. await child.waitForExit()
  141. } catch (error: unknown) {
  142. failures.push(thrown(error))
  143. }
  144. }
  145. try {
  146. await child.done
  147. } catch (error: unknown) {
  148. failures.push(thrown(error))
  149. }
  150. const firstFailure = failures[0]
  151. if (failures.length === 1 && firstFailure !== undefined) throw firstFailure
  152. if (failures.length > 1) {
  153. throw new AggregateError(
  154. failures,
  155. 'subagent-claude-code: query and process cleanup failed',
  156. )
  157. }
  158. }
  159. /**
  160. * Build the fixed official SDK options for one one-shot provider run.
  161. * @param spec - Workspace, environment, process service, and disposal policy.
  162. * @param controller - per-run cancellation owner.
  163. * @param capture - receives the real managed child synchronously from the SDK hook.
  164. * @returns options that inherit native settings while disabling persistence and user questions.
  165. */
  166. export function claudeQueryOptions(
  167. spec: ClaudeCodeRunSpec,
  168. controller: AbortController,
  169. capture: (child: SubprocessHandle) => void,
  170. ): Options {
  171. return {
  172. abortController: controller,
  173. cwd: spec.cwd,
  174. pathToClaudeCodeExecutable: spec.executable,
  175. env: { ...scrubbedParentEnv(), ...spec.env },
  176. persistSession: false,
  177. disallowedTools: ['AskUserQuestion'],
  178. spawnClaudeCodeProcess: (options: SpawnOptions) => {
  179. const child = spec.spawn(claudeSpawnSpec(options, spec.disposeGraceMs))
  180. capture(child)
  181. return new ManagedClaudeCodeProcess(child)
  182. },
  183. }
  184. }
  185. /**
  186. * Start one official Claude Agent SDK query and publish its one-shot run.
  187. * @param request - resolved shared subagent request.
  188. * @param spec - Workspace, environment, process service, and diagnostic policy.
  189. * @returns the published run after both Query and real CLI handle exist.
  190. */
  191. export async function startClaudeCodeRun(
  192. request: SubagentStartRequest,
  193. spec: ClaudeCodeRunSpec,
  194. ): Promise<SubagentRun> {
  195. const prompt = textTask(request.prompt)
  196. if (request.signal.aborted) {
  197. throw new Error('subagent-claude-code: request was aborted before SDK startup')
  198. }
  199. const controller = new AbortController()
  200. const requestCancel = (): void => {
  201. if (!controller.signal.aborted) {
  202. controller.abort(new Error('subagent-claude-code: run cancelled locally'))
  203. }
  204. }
  205. const onAbort = (): void => { requestCancel() }
  206. request.signal.addEventListener('abort', onAbort, { once: true })
  207. let child: SubprocessHandle | undefined
  208. let query: Query | undefined
  209. try {
  210. query = officialQuery({
  211. prompt,
  212. options: claudeQueryOptions(spec, controller, (captured) => {
  213. child = captured
  214. }),
  215. })
  216. if (child === undefined || child.pid <= 0) {
  217. throw new Error(
  218. 'subagent-claude-code: official SDK did not publish a controllable Claude Code process',
  219. )
  220. }
  221. if (controller.signal.aborted) {
  222. throw new Error('subagent-claude-code: request was aborted before SDK startup')
  223. }
  224. } catch (error: unknown) {
  225. request.signal.removeEventListener('abort', onAbort)
  226. const cancelledBeforeCleanup = controller.signal.aborted
  227. requestCancel()
  228. if (child !== undefined) {
  229. try {
  230. await disposeClaudeCodeChild(query, child)
  231. } catch (disposeError: unknown) {
  232. throw new AggregateError(
  233. [thrown(error), thrown(disposeError)],
  234. 'subagent-claude-code: startup failed and CLI cleanup also failed',
  235. )
  236. }
  237. } else if (query !== undefined) {
  238. try {
  239. query.close()
  240. } catch (disposeError: unknown) {
  241. throw new AggregateError(
  242. [thrown(error), thrown(disposeError)],
  243. 'subagent-claude-code: startup failed and query cleanup also failed',
  244. )
  245. }
  246. }
  247. // oxlint-disable-next-line typescript/no-unnecessary-condition -- the request can abort while process cleanup is awaited.
  248. if (cancelledBeforeCleanup || request.signal.aborted) {
  249. throw new Error('subagent-claude-code: request was aborted before SDK startup')
  250. }
  251. throw thrown(error)
  252. }
  253. const publishedQuery = query
  254. const publishedChild = child
  255. const result = settleRunResult({
  256. attempt: () => consumeClaudeQuery(publishedQuery),
  257. collectOutput: () => [],
  258. cancelled: () => controller.signal.aborted,
  259. onError: spec.onError,
  260. signal: request.signal,
  261. onAbort,
  262. })
  263. return subprocessRunHandle({
  264. id: SessionId(randomUUID()),
  265. result,
  266. signal: request.signal,
  267. onAbort,
  268. requestCancel,
  269. teardown: () => disposeClaudeCodeChild(
  270. publishedQuery,
  271. publishedChild,
  272. ),
  273. })
  274. }