structured.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312
  1. /**
  2. * Structured-output support for the in-process subagent backends: the mechanism
  3. * behind `SubagentStartRequest.outputSchema` for children that run as agents on
  4. * the same context.
  5. *
  6. * The model-facing surface is one globally registered `structured_output` tool
  7. * whose REGISTERED parameters are a placeholder — the real schema is per run.
  8. * Because the tool registry and prompt assembly are context-global while
  9. * schemas differ per child (two concurrent structured runs may carry different
  10. * schemas), per-agent shaping happens on the `system-prompt/assemble`
  11. * waterfall with a `prepend: true` listener that post-processes `await next()`
  12. * — FINAL-ASSEMBLY enforcement: whatever downstream listeners mutated or
  13. * replaced, the assembly the loop renders never carries `structured_output`
  14. * for an agent without a structured run, and for one that has it always
  15. * carries the run's OWN schema plus a trailing
  16. * {@link STRUCTURED_OUTPUT_INSTRUCTION} section (the demand travels with the
  17. * tool). The loop logs what the assembly produced as the request header, so
  18. * the injection is a reconstructable fact of the session log, never a
  19. * wire-only mutation (the reconstructability RFC).
  20. * (Cooperative mutate-then-`next()` would not survive a downstream listener
  21. * returning a replacement assembly — see the waterfall composition caveat in
  22. * docs/architecture.md.)
  23. *
  24. * FIXME: the whole enforcement dance above exists because the tool registry
  25. * and prompt assembly are context-global. If they become per-agent or
  26. * per-session scoped, a structured run just registers its own schema'd tool on
  27. * the child's scope and this module reduces to the capture tool plus the
  28. * turn-stop — no placeholder, no final-assembly swap, no strip-for-everyone-
  29. * else, no global-registration lifetime dance.
  30. *
  31. * A companion `agent/turn-continuation` listener stops a child's turn once its
  32. * output is captured — without it, the loop's default "had tool calls ⇒
  33. * continue" buys a wasted extra model step per structured child. It is also
  34. * `prepend: true`: the veto must run before any earlier-registered listener
  35. * that could short-circuit the chain into a forced continue. A third listener
  36. * closes the within-step window the continuation veto cannot: a
  37. * `tools/pre-execute` deny for any call arriving after the agent's capture, so
  38. * a response that lists `structured_output` before further tool calls cannot
  39. * run side effects after the final answer was accepted. A fourth,
  40. * `tools/post-execute`, is the capture COMMIT: the tool body only stages the
  41. * validated value, and it becomes the run's captured result only when the
  42. * final post-execute decision accepts the call — a blocking hook downstream
  43. * yields `isError` in the log, and the run must not report success for it.
  44. *
  45. * Lifetime is refcounted by structured RUNS: each acquires from start to
  46. * settle, so the registrations exist exactly while at least one structured
  47. * child is live — a plain deployment that never passes `outputSchema` carries
  48. * no always-on global state, and a backend hot-reload mid-run cannot
  49. * unregister the capture tool out from under a live child (the run holds its
  50. * own acquisition). Registrations land on the ROOT context and the refcount
  51. * disposes them when the last run settles; the next structured run
  52. * re-registers them.
  53. *
  54. * @module @deepseek-ai/dsh-subagent-inprocess/structured
  55. */
  56. import type { Context } from 'cordis'
  57. import type { Agent } from '@deepseek-ai/dsh-agent'
  58. import type { ContentBlock, ToolSchema } from '@deepseek-ai/dsh-llm'
  59. import type { ContinuationDecision } from '@deepseek-ai/dsh-agent'
  60. import type { AssembleContext, PromptAssembly } from '@deepseek-ai/dsh-system-prompt'
  61. import type { PostToolDecision, PreToolDecision, ToolExecution, ToolExecutionResult } from '@deepseek-ai/dsh-tools'
  62. import { ToolArgsError, validateStructuredValue, type StructuredOutputSchema } from '@deepseek-ai/dsh-tools'
  63. /** The model-facing tool name a structured child must call to finish. */
  64. export const STRUCTURED_OUTPUT_TOOL = 'structured_output'
  65. /**
  66. * The instruction the assembly listener appends to a structured child's
  67. * system prompt as a trailing section on every assembly. Per-assembly state,
  68. * NOT agent prompt state: `AgentOptions` has no prompt field (the persona is
  69. * deployment config on the system-prompt plugin), so the same final-assembly
  70. * enforcement that injects the schema'd tool carries the instruction that
  71. * demands calling it.
  72. */
  73. export const STRUCTURED_OUTPUT_INSTRUCTION
  74. = 'When you have your final answer, you MUST report it by calling the '
  75. + `\`${STRUCTURED_OUTPUT_TOOL}\` tool with arguments matching its parameter schema exactly. `
  76. + 'Do not finish with a plain text answer: only the tool call counts as your result.'
  77. /** One structured run's state: the schema to enforce and the captured value, once recorded. */
  78. interface RunState {
  79. readonly schema: StructuredOutputSchema
  80. /**
  81. * A validated value awaiting the post-execute verdict on ITS OWN call. Set
  82. * by the capture tool's body, promoted to {@link RunState.captured} only
  83. * when the final `tools/post-execute` decision accepts the call — a
  84. * downstream block turns the logged result into `isError`, and a value
  85. * committed at body time would let the run report success for a call the
  86. * model saw fail.
  87. */
  88. pending?: { value: unknown }
  89. captured?: { value: unknown }
  90. }
  91. /** The per-root-context runtime: run states plus the shared registrations. */
  92. interface StructuredRuntime {
  93. refs: number
  94. readonly states: WeakMap<Agent, RunState>
  95. readonly disposers: (() => void)[]
  96. }
  97. /** One root context ⇒ one runtime (multi-app test isolation). */
  98. const runtimes = new WeakMap<Context, StructuredRuntime>()
  99. /**
  100. * One holder's handle on the shared structured runtime. `release()` is
  101. * idempotent per acquisition; the runtime's registrations are disposed when the
  102. * LAST holder (backend plugin or live run) releases.
  103. */
  104. export interface StructuredAcquisition {
  105. /** Enforce `schema` on `agent`'s requests and start capturing its `structured_output` call. */
  106. attach(agent: Agent, schema: StructuredOutputSchema): void
  107. /** The captured value, once the child called the tool with valid arguments. */
  108. captured(agent: Agent): { value: unknown } | undefined
  109. /** Stop enforcing/capturing for `agent` (WeakMap-backed; safe to call twice). */
  110. detach(agent: Agent): void
  111. /** Drop this holder's reference (idempotent); the last release unregisters everything. */
  112. release(): void
  113. }
  114. /**
  115. * Acquire the per-root-context structured runtime, registering the capture tool
  116. * and the runtime's listeners on the FIRST acquisition. See the module doc
  117. * for the enforcement and lifetime design.
  118. * @param ctx - any context of the app; the runtime keys off `ctx.root`.
  119. * @returns this holder's handle (attach/captured/detach + idempotent release).
  120. */
  121. export function acquireStructuredRuntime(ctx: Context): StructuredAcquisition {
  122. const root: Context = ctx.root
  123. let runtime = runtimes.get(root)
  124. if (!runtime) {
  125. runtime = { refs: 0, states: new WeakMap(), disposers: [] }
  126. runtimes.set(root, runtime)
  127. registerRuntime(root, runtime)
  128. }
  129. runtime.refs += 1
  130. let released = false
  131. return {
  132. attach(agent: Agent, schema: StructuredOutputSchema): void {
  133. runtime.states.set(agent, { schema })
  134. },
  135. captured(agent: Agent): { value: unknown } | undefined {
  136. return runtime.states.get(agent)?.captured
  137. },
  138. detach(agent: Agent): void {
  139. runtime.states.delete(agent)
  140. },
  141. release(): void {
  142. if (released) return
  143. released = true
  144. runtime.refs -= 1
  145. if (runtime.refs > 0) return
  146. runtimes.delete(root)
  147. for (const dispose of runtime.disposers.splice(0)) dispose()
  148. },
  149. }
  150. }
  151. /** Register the capture tool + the two listeners on the root context (first acquire). */
  152. function registerRuntime(root: Context, runtime: StructuredRuntime): void {
  153. // The registered parameters are a PLACEHOLDER: the request listener below
  154. // swaps in the run's real schema per child, and strips the tool entirely for
  155. // every agent without a structured run — so this shape is never model-visible.
  156. //
  157. // Registration does NOT ride on the acquiring backend's plugin-level
  158. // `inject`: a backend that waited on `tools` would apply later than it did
  159. // before this module existed, shifting when its PROVIDER registers — and the
  160. // delegation tool mirrors provider lifecycle, so that shift would reorder
  161. // the model-visible tool list of every existing prompt. Instead the capture
  162. // tool registers synchronously when `tools` is already live (the common
  163. // case), and through a scoped inject fiber when the Loader happens to start
  164. // the backend first. Either way the registration lands on root and is
  165. // disposed by the runtime's refcount; disposing the fiber also covers the
  166. // never-activated case.
  167. let disposeTool: (() => void) | undefined
  168. const registerCapture = (tools: Context['tools']): void => {
  169. disposeTool = tools.register({
  170. name: STRUCTURED_OUTPUT_TOOL,
  171. description:
  172. 'Report your final structured result. Call this exactly once, when your answer is complete; '
  173. + 'the arguments must match this tool\'s parameter schema exactly.',
  174. parameters: { type: 'object', properties: {} },
  175. execute(args: unknown, exec: ToolExecution): Promise<ContentBlock[]> {
  176. const state = exec.agent ? runtime.states.get(exec.agent) : undefined
  177. if (!state) {
  178. // Reachable only if a non-structured agent somehow calls the tool (the
  179. // request listener strips it, so the model never sees it) — fail loud
  180. // rather than capture into nowhere.
  181. throw new Error(`${STRUCTURED_OUTPUT_TOOL} is only available to subagents started with an output schema`)
  182. }
  183. const violations = validateStructuredValue(state.schema, args)
  184. // ToolArgsError → isError result with INVALID_ARGS: the model retries
  185. // within the same turn, exactly like a schema-validated defineTool call.
  186. if (violations.length > 0) throw new ToolArgsError(violations)
  187. // Two-phase commit: the body only STAGES the value; the post-execute
  188. // listener below promotes it once the final decision accepts the call.
  189. state.pending = { value: args }
  190. return Promise.resolve([{ type: 'text', text: 'Structured output recorded.' }])
  191. },
  192. })
  193. }
  194. const liveTools = root.get('tools')
  195. const toolsFiber = liveTools ? undefined : root.inject(['tools'], (childCtx: Context) => {
  196. registerCapture(childCtx.root.tools)
  197. })
  198. if (liveTools) registerCapture(liveTools)
  199. runtime.disposers.push(() => {
  200. disposeTool?.()
  201. void toolsFiber?.dispose()
  202. })
  203. // FINAL-ASSEMBLY enforcement (prepend: true = first registered = OUTERMOST
  204. // wrapper): post-process whatever the downstream listeners and the registry
  205. // produced, so a downstream listener returning a replacement assembly cannot
  206. // leak the tool to other agents or erase the child's schema. The loop logs
  207. // the rendered assembly as the step's request header, so the swap is
  208. // reconstructable log state, never a wire-only mutation.
  209. runtime.disposers.push(root.on('system-prompt/assemble', async function (
  210. this: unknown, _assembly: PromptAssembly, context: AssembleContext, next: () => Promise<PromptAssembly>,
  211. ): Promise<PromptAssembly> {
  212. const final = await next()
  213. const state = context.agent ? runtime.states.get(context.agent) : undefined
  214. if (state) {
  215. const schemaEntry: ToolSchema = {
  216. name: STRUCTURED_OUTPUT_TOOL,
  217. description:
  218. 'Report your final structured result. Call this exactly once, when your answer is complete; '
  219. + 'the arguments must match this tool\'s parameter schema exactly.',
  220. // ToolSchema.parameters is the wire-level JSON Schema object; the
  221. // asserted subset type is structurally exactly that.
  222. parameters: state.schema as unknown as Record<string, unknown>,
  223. }
  224. final.tools = [...final.tools.filter(tool => tool.name !== STRUCTURED_OUTPUT_TOOL), schemaEntry]
  225. // The demand travels WITH the tool: a trailing section in the
  226. // tool-guidance order band, appended after next() so it renders last
  227. // (renderPrompt joins in array order).
  228. final.sections = [...final.sections, { name: `tool:${STRUCTURED_OUTPUT_TOOL}`, order: 190, text: STRUCTURED_OUTPUT_INSTRUCTION }]
  229. return final
  230. }
  231. // No structured run: strip the placeholder so it is never model-visible.
  232. // An empty tools array canonicalizes to an absent header/wire field
  233. // (canonicalHeader pins empty ≡ absent), so no re-shaping is needed here.
  234. final.tools = final.tools.filter(tool => tool.name !== STRUCTURED_OUTPUT_TOOL)
  235. return final
  236. }, { prepend: true }))
  237. // Stop a structured child's turn once its output is captured: the default
  238. // "had tool calls ⇒ continue" would otherwise buy a wasted extra model step
  239. // after every successful capture. `prepend: true` puts the veto OUTERMOST —
  240. // an earlier-registered listener that short-circuits the chain (a goal-style
  241. // force-continue returning without `next()`) would otherwise decide the turn
  242. // before this listener ever ran, and no downstream decision may resurrect a
  243. // structured turn that is already finished.
  244. runtime.disposers.push(root.on('agent/turn-continuation', function (
  245. this: unknown, agent: Agent, _turn: number, _decision: ContinuationDecision, next: () => Promise<ContinuationDecision>,
  246. ): Promise<ContinuationDecision> {
  247. if (runtime.states.get(agent)?.captured) return Promise.resolve({ action: 'stop' })
  248. return next()
  249. }, { prepend: true }))
  250. // The capture COMMIT: promote the staged value only when the final
  251. // post-execute decision accepts the call. The capture tool's body cannot
  252. // decide — `tools/post-execute` runs after it, and a blocking listener (a
  253. // PostToolUse hook) turns the logged result into `isError` feedback; a value
  254. // committed at body time would make readResult report `structured` success
  255. // for a call whose result the model and session log saw fail. `prepend:
  256. // true` = outermost at registration time, so `await next()` returns the
  257. // COMPOSED downstream decision — the same final verdict the registry maps
  258. // onto the result. (A later-registered outer listener that blocks without
  259. // delegating skips this commit entirely: the staged value is dropped and the
  260. // run errors — failure-safe in the same direction.) The staging slot clears
  261. // on every path, including a rejecting downstream listener.
  262. runtime.disposers.push(root.on('tools/post-execute', async function (
  263. this: unknown, exec: ToolExecution, _result: ToolExecutionResult, next: () => Promise<PostToolDecision>,
  264. ): Promise<PostToolDecision> {
  265. const state = exec.agent ? runtime.states.get(exec.agent) : undefined
  266. if (!state || exec.name !== STRUCTURED_OUTPUT_TOOL || state.pending === undefined) return next()
  267. const pending = state.pending
  268. try {
  269. const decision = await next()
  270. if (decision.kind === 'accept') state.captured = pending
  271. return decision
  272. } finally {
  273. delete state.pending
  274. }
  275. }, { prepend: true }))
  276. // Terminal means terminal WITHIN the step, not only at its end: the
  277. // turn-continuation veto above runs after every call in the current model
  278. // response has executed, so a response that puts `structured_output` before
  279. // further tool calls would still perform those side effects after the final
  280. // answer was accepted. Deny every later call for a captured agent at the
  281. // allow/deny gate — dispatch is skipped and the model sees an `isError`
  282. // result naming the contract. Calls that PRECEDE the capture in the same
  283. // response ran before `captured` was set and are untouched; a second
  284. // `structured_output` is denied like any other call. `prepend: true` for the
  285. // same reason as the continuation veto: no earlier-registered allow may
  286. // short-circuit past the terminal contract.
  287. runtime.disposers.push(root.on('tools/pre-execute', function (
  288. this: unknown, exec: ToolExecution, next: () => Promise<PreToolDecision>,
  289. ): Promise<PreToolDecision> {
  290. if (exec.agent && runtime.states.get(exec.agent)?.captured) {
  291. return Promise.resolve({
  292. kind: 'deny',
  293. reason: `structured output already recorded: the run is complete, so \`${exec.name}\` is not executed`,
  294. })
  295. }
  296. return next()
  297. }, { prepend: true }))
  298. }