返回源码地图

packages/workflow/tool-workflow/src/index.ts

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

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

1/**
2 * The model-facing `workflow` tool: run a JavaScript orchestration script that fans out
3 * subagents, and return the script's final value. It owns the model-facing schema and run lifecycle; script
4 * parsing, execution, caps, and cancellation live behind `ctx.workflowEngine`
5 * (`@deepseek-ai/dsh-workflow`), so a hardened engine swaps in without touching what the model
6 * sees. Foreground execution awaits `run.result` and always disposes the run; non-completed reasons
7 * become tool errors. `run_in_background: true` instead registers the run as an owned `ctx.jobs` job
8 * and returns its id immediately — the job's output ring streams live progress, and the run's value
9 * arrives with the job's completion notice. Presentation is an args-only generic card
10 * titled from `meta.name`. Explicit-ask usage guidance is registered as the tool's own prompt
11 * section rather than deployment persona prose.
12 * @module @deepseek-ai/dsh-tool-workflow
13 */
14
15import type { Context } from '@deepseek-ai/cordis'
16import z from '@deepseek-ai/schemastery'
17import { defineTool } from '@deepseek-ai/dsh-tools'
18import type { ToolCallView, ToolResultView } from '@deepseek-ai/dsh-tools'
19import type { Agent } from '@deepseek-ai/dsh-agent'
20import type { ContentBlock } from '@deepseek-ai/dsh-llm'
21import type { JobId, JobOutcome } from '@deepseek-ai/dsh-jobs'
22import type { Session, SessionEventMap } from '@deepseek-ai/dsh-session'
23import type { JsonValue } from '@deepseek-ai/dsh-util-values'
24import type {
25 WorkflowResult, WorkflowRun, WorkflowRunId, WorkflowStopReason,
26} from '@deepseek-ai/dsh-workflow'
27import { createWorkflowRecordMirror } from './record.ts'
28import type { WorkflowRecordMirror } from './record.ts'
29import type {
30 ToolWorkflowAgentEndData, ToolWorkflowAgentStartData,
31 ToolWorkflowRunEndData, ToolWorkflowRunStartData,
32} from './types.ts'
33
34declare module '@deepseek-ai/dsh-jobs' {
35 interface JobKindMap {
36 workflow: 'workflow'
37 }
38}
39
40export const name = 'tool-workflow'
41export const inject = ['tools', 'workflowEngine', 'systemPrompt']
42
43/** Config: the model-facing tool name plus result rendering caps. */
44export interface Config {
45 /** The model-facing tool name to register (default `workflow`). */
46 toolName?: string
47 /** Rendered-result ceiling, in characters: a longer JSON value is truncated with a notice (default 50000). */
48 maxResultChars?: number
49 /**
50 * Expose `run_in_background` (default true); disabled calls are also
51 * rejected. A background run needs a live `ctx.jobs` registry with a
52 * controller serving the caller (`dsh-jobs-local` plus `dsh-tool-jobs` in
53 * the shipped composition); without one the call fails with the missing
54 * piece named.
55 */
56 enableRunInBackground?: boolean
57}
58
59export const Config: z<Config> = z.object({
60 toolName: z.string().default('workflow'),
61 maxResultChars: z.natural().min(1).default(50_000),
62 enableRunInBackground: z.boolean().default(true),
63})
64
65type ResolvedConfig = Required<Config>
66
67interface WorkflowRecorder {
68 start(session: Session, run: WorkflowRun): void
69 finish(runId: WorkflowRunId, stopReason: WorkflowStopReason): void
70 abandon(runId: WorkflowRunId): void
71}
72
73interface ToolWorkflowRecordEventMap {
74 'tool-workflow/run-start': ToolWorkflowRunStartData
75 'tool-workflow/agent-start': ToolWorkflowAgentStartData
76 'tool-workflow/agent-end': ToolWorkflowAgentEndData
77 'tool-workflow/run-end': ToolWorkflowRunEndData
78}
79
80/** Render a contained recording failure without trusting the thrown value. */
81function renderRecordingError(error: unknown): string {
82 try {
83 return String(error)
84 } catch {
85 return '[unrenderable thrown value]'
86 }
87}
88
89/**
90 * Project active top-level workflow runs into their parent Sessions without
91 * letting recording failure affect tool execution.
92 */
93function createWorkflowRecorder(ctx: Context): WorkflowRecorder {
94 const active = new Map<WorkflowRunId, Session>()
95 const append = <Type extends keyof ToolWorkflowRecordEventMap>(
96 session: Session,
97 type: Type,
98 data: SessionEventMap[Type],
99 ): boolean => {
100 // These four package-owned events are all log-only. Narrowing the generic
101 // append face here discharges Session.append's conditional options tuple.
102 const appendRecord = session.append.bind(session) as <Event extends keyof ToolWorkflowRecordEventMap>(
103 event: Event,
104 value: SessionEventMap[Event],
105 ) => void
106 try {
107 appendRecord(type, data)
108 return true
109 } catch (error: unknown) {
110 ctx.logger.warn(`tool-workflow: disabled durable record after ${type} append failed: ${renderRecordingError(error)}`)
111 return false
112 }
113 }
114
115 ctx.on('workflow/agent-start', (info, agent) => {
116 const session = active.get(info.id)
117 if (session === undefined) return
118 const data: ToolWorkflowAgentStartData = {
119 runId: info.id,
120 seq: agent.seq,
121 label: agent.label,
122 ...agent.phase === undefined ? {} : { phase: agent.phase },
123 childId: agent.childId,
124 }
125 if (!append(session, 'tool-workflow/agent-start', data)) active.delete(info.id)
126 })
127 ctx.on('workflow/agent-end', (info, agent) => {
128 const session = active.get(info.id)
129 if (session === undefined) return
130 const data: ToolWorkflowAgentEndData = {
131 runId: info.id,
132 seq: agent.seq,
133 outcome: agent.outcome,
134 }
135 if (!append(session, 'tool-workflow/agent-end', data)) active.delete(info.id)
136 })
137
138 return {
139 start(session, run) {
140 if (append(session, 'tool-workflow/run-start', { runId: run.id, name: run.meta.name })) {
141 active.set(run.id, session)
142 }
143 },
144 finish(runId, stopReason) {
145 const session = active.get(runId)
146 if (session !== undefined) append(session, 'tool-workflow/run-end', { runId, stopReason })
147 active.delete(runId)
148 },
149 abandon: (runId) => { active.delete(runId) },
150 }
151}
152
153/**
154 * The script-authoring contract, embedded in the tool description: the hooks,
155 * their exact semantics, and the supported schema subset. Parameter-level
156 * rules live in the parameter descriptions.
157 */
158const DESCRIPTION = `Run a JavaScript workflow script that orchestrates subagents at scale. Use this for work that fans out across many independent pieces — an audit over many files, a migration, multi-angle research, adversarial verification of findings — where you write the orchestration as a script instead of delegating turn by turn.
159
160Script-body hooks:
161- \`agent(prompt, opts?): Promise<any>\` — run one subagent to completion. Without \`opts.schema\` it resolves to the child's final text; with \`opts.schema\` (an object-rooted JSON Schema using ONLY type/properties/required/additionalProperties/items/enum/const/oneOf) it resolves to the validated object. Resolves \`null\` when the child fails (filter with \`.filter(Boolean)\`). Other opts: \`label\` (display), \`phase\` (progress group), and independent \`provider\`/\`model\` LLM target overrides.
162- \`pipeline(items, ...stages): Promise<any[]>\` — run each item through the stages independently with NO barrier between stages (prefer this for multi-stage work). Each stage receives \`(prev, item, index)\`. A stage throw drops that ITEM to \`null\` and skips its remaining stages.
163- \`parallel(thunks): Promise<any[]>\` — run zero-argument functions concurrently and await ALL of them (a barrier; use only when a stage genuinely needs every prior result together). A throwing thunk resolves to \`null\`.
164- \`phase(title)\` — start a progress phase; \`log(message)\` — narrate progress; \`args\` — the tool call's \`args\` input, verbatim.
165
166Misused hooks (bad arguments, unknown options, unsupported schemas, tripped caps) end the whole script instead of producing \`null\`. The script has no filesystem, network, timer, or Node.js APIs; the agents do the work.`
167
168type WorkflowCallArgs = {
169 script: string
170 meta: {
171 name: string
172 description: string
173 whenToUse?: string
174 phases?: { title: string; detail?: string; provider?: string; model?: string }[]
175 }
176 args?: Record<string, unknown>
177 run_in_background?: boolean
178}
179
180/** The pending-state card: a generic card titled by the workflow's meta name. */
181function presentWorkflowCall(args: WorkflowCallArgs): ToolCallView {
182 return {
183 card: 'generic',
184 title: `workflow: ${args.meta.name}`,
185 rawInput: args.script,
186 }
187}
188
189/** The completed-state card: keep the pending title; render the result content as-is. */
190function presentWorkflowResult(args: WorkflowCallArgs, result: { content: ContentBlock[]; isError: boolean }): ToolResultView {
191 void args
192 void result
193 return { card: 'generic' }
194}
195
196/** A non-`completed` stop reason means the script did not finish cleanly. */
197function stopReasonError(result: WorkflowResult): string | undefined {
198 switch (result.stopReason) {
199 case 'completed':
200 return undefined
201 case 'cancelled':
202 return `workflow run was cancelled${result.error !== undefined ? ` (${result.error})` : ''}`
203 case 'error':
204 return `workflow run failed: ${result.error ?? 'unknown error'}`
205 /* v8 ignore start -- defensive: WorkflowStopReason is a closed union, exhaustive by construction; a future variant fails here loudly */
206 default:
207 return `workflow run ended abnormally (${String(result.stopReason satisfies never)})`
208 /* v8 ignore stop */
209 }
210}
211
212/**
213 * Map a settled background run onto the job outcome vocabulary. A completed
214 * run carries the rendered return value as the job's result; a
215 * cancelled run leaves the detail to the registry's kill-reason merge (the
216 * cancel reason it forwarded is the same string); an errored run fails with
217 * the script's failure message.
218 */
219function jobOutcomeOf(result: WorkflowResult, name: string, maxChars: number): JobOutcome {
220 switch (result.stopReason) {
221 case 'completed':
222 return {
223 status: 'completed',
224 detail: `${result.agentsStarted} agent${result.agentsStarted === 1 ? '' : 's'}`,
225 result: renderResult(name, result.agentsStarted, result.value as JsonValue, maxChars),
226 }
227 case 'cancelled':
228 return { status: 'killed' }
229 case 'error':
230 return { status: 'failed', detail: result.error ?? 'unknown error' }
231 /* v8 ignore start -- defensive: WorkflowStopReason is a closed union, exhaustive by construction; a future variant fails here loudly */
232 default:
233 return { status: 'failed', detail: `workflow run ended abnormally (${String(result.stopReason satisfies never)})` }
234 /* v8 ignore stop */
235 }
236}
237
238/** Render the run's outcome text: the meta name, agent count, and the JSON value (capped). */
239function renderResult(name: string, agentsStarted: number, value: JsonValue, maxChars: number): string {
240 // The engine returns JSON data (null for a valueless script), so stringify never yields undefined.
241 const rendered = JSON.stringify(value, null, 2)
242 const clipped = rendered.length > maxChars
243 ? `${rendered.slice(0, maxChars)}\n… [truncated: ${rendered.length - maxChars} more characters]`
244 : rendered
245 return `workflow "${name}" completed (${agentsStarted} agent${agentsStarted === 1 ? '' : 's'}).\nReturn value:\n${clipped}`
246}
247
248/**
249 * Register a background run as an owned job. The engine run is started
250 * inside the job starter with no tool-step signal — the run belongs to the
251 * job, so a registry kill or owner teardown is what cancels it — and its
252 * settlement is the job's settlement: dispose, stop the mirrors, then map the
253 * stop reason onto the job outcome (a completed run's rendered return value
254 * rides `result` to the model's first read after settlement).
255 * @param ctx - plugin context (engine, optional jobs registry, logger).
256 * @param args - the validated tool call.
257 * @param parent - the calling agent; owns the job.
258 * @param recordsRun - whether this top-level call records durable run events.
259 * @param deps - the tool's recorder/mirror taps and the render cap.
260 * @returns the background result for the tool's output schema.
261 */
262function startBackgroundRun(
263 ctx: Context,
264 args: WorkflowCallArgs,
265 parent: Agent,
266 recordsRun: boolean,
267 deps: { recorder: WorkflowRecorder; mirror: WorkflowRecordMirror; maxResultChars: number },
268): { kind: 'background'; jobId: JobId; runId: WorkflowRunId } {
269 const jobs = ctx.get('jobs')
270 if (jobs === undefined) {
271 throw new Error('background jobs unavailable: load @deepseek-ai/dsh-jobs and @deepseek-ai/dsh-tool-jobs')
272 }
273 let run!: WorkflowRun
274 const jobId = jobs.start({
275 kind: 'workflow',
276 label: args.meta.name,
277 owner: parent.id,
278 run: (job) => {
279 // A synchronous engine rejection (META_INVALID/SCRIPT_PARSE) propagates
280 // out of the starter, so the registry registers nothing and the model
281 // sees the violation list as an ordinary tool error.
282 run = ctx.workflowEngine.start({
283 script: args.script,
284 meta: args.meta,
285 ...args.args !== undefined ? { args: args.args } : {},
286 parent,
287 })
288 deps.mirror.start(run.id, job)
289 if (recordsRun) deps.recorder.start(parent.session, run)
290 const done = run.result.then(async (result): Promise<JobOutcome> => {
291 try {
292 // Keep member listeners alive through disposal: an engine may
293 // synthesize cancelled member endings while reaching quiescence.
294 await run.dispose()
295 } catch (error: unknown) {
296 // done must not reject; a failed disposal still has a settled result to report.
297 ctx.logger.warn(`background workflow run ${run.id} dispose failed: ${String(error)}`)
298 }
299 deps.mirror.stop(run.id)
300 if (recordsRun) {
301 deps.recorder.finish(run.id, result.stopReason)
302 deps.recorder.abandon(run.id)
303 }
304 return jobOutcomeOf(result, args.meta.name, deps.maxResultChars)
305 })
306 return {
307 cancel: (reason?: string) => { run.cancel(reason ?? 'background workflow job killed') },
308 done,
309 }
310 },
311 })
312 return { kind: 'background' as const, jobId, runId: run.id }
313}
314
315export function apply(ctx: Context, config: Config): void {
316 // schemastery (the exported Config schema) has already filled the defaulted
317 // fields; the assertion records that resolution, not a hidden fallback.
318 const { toolName, maxResultChars, enableRunInBackground } = config as ResolvedConfig
319 const recorder = createWorkflowRecorder(ctx)
320 const mirror = createWorkflowRecordMirror(ctx)
321 // Usage policy ships with the tool (the master convention: tool guidance
322 // lives in tool plugins as prompt sections, not in the deployment persona).
323 ctx.systemPrompt.section({
324 name: `tool:${toolName}`,
325 order: ctx.systemPrompt.getSectionOrder('TOOL_WORKFLOW'),
326 text: `Use the ${toolName} tool ONLY when the user explicitly asks for a workflow or for large multi-agent orchestration: you write a JavaScript script (the tool description documents the exact format) that fans work out across many subagents with phases and structured results. For one or two delegations, prefer plain subagent calls.`,
327 })
328 ctx.tools.register(defineTool({
329 name: toolName,
330 description: DESCRIPTION,
331 parameters: {
332 script: {
333 type: 'string',
334 required: true,
335 description: 'The plain JavaScript body, not TypeScript and without an `export const meta` statement; top-level await is allowed. '
336 + 'End with `return <value>`; the JSON-serializable value is this tool\'s result.',
337 },
338 meta: {
339 type: 'object',
340 additionalProperties: true,
341 required: true,
342 description: 'The workflow identity as plain JSON, not code.',
343 properties: {
344 name: { type: 'string', required: true, description: 'Short kebab-case workflow name.' },
345 description: { type: 'string', required: true, description: 'One-line description of what the workflow does.' },
346 whenToUse: { type: 'string', description: 'Optional guidance on when this workflow applies.' },
347 phases: {
348 type: 'array',
349 description: 'Optional phase declarations matched by phase() calls.',
350 items: {
351 type: 'object',
352 additionalProperties: true,
353 properties: {
354 title: { type: 'string', required: true, description: 'The phase title phase() calls match by exact string.' },
355 detail: { type: 'string', description: 'Optional one-line description of the phase.' },
356 provider: { type: 'string', description: 'Optional provider override this phase is expected to use.' },
357 model: { type: 'string', description: 'Optional model override this phase is expected to use.' },
358 },
359 },
360 },
361 },
362 },
363 args: {
364 type: 'object',
365 additionalProperties: true,
366 description: 'Optional JSON input exposed to the script as the `args` global (wrap a bare list as a field, e.g. {"files": [...]}).',
367 },
368 ...enableRunInBackground ? {
369 run_in_background: {
370 type: 'boolean' as const,
371 description: 'Run as a background job: return a job id immediately instead of waiting; the return value arrives with the completion notice.',
372 },
373 } : {},
374 },
375 output: {
376 schema: {
377 oneOf: [
378 {
379 type: 'object',
380 additionalProperties: false,
381 properties: {
382 kind: { type: 'string', required: true, const: 'background' },
383 jobId: { type: 'string', required: true },
384 runId: { type: 'string', required: true },
385 },
386 },
387 {
388 type: 'object',
389 additionalProperties: false,
390 properties: {
391 kind: { type: 'string', required: true, const: 'foreground' },
392 runId: { type: 'string', required: true },
393 agentsStarted: { type: 'integer', required: true },
394 result: { type: 'json', required: true },
395 },
396 },
397 ],
398 },
399 render: (args, value) => [{
400 type: 'text',
401 text: value.kind === 'background'
402 ? `workflow "${args.meta.name}" started in the background as job ${value.jobId}. Its return value arrives with the completion notice; check on it with job_output, stop it with job_kill.`
403 : renderResult(args.meta.name, value.agentsStarted, value.result, maxResultChars),
404 }],
405 },
406 async execute(args, exec) {
407 const parent = exec.agent
408 if (!parent) {
409 // The loop sets `exec.agent` for every model-driven call; its absence
410 // means a non-agent caller invoked the tool directly, which has no
411 // parent to attribute the children to. Fail loud rather than guess.
412 throw new Error('workflow tool requires a calling agent (exec.agent was undefined)')
413 }
414 if (args.run_in_background === true) {
415 if (!enableRunInBackground) {
416 throw new Error('run_in_background is disabled for this tool')
417 }
418 // No pre-abort check here, unlike bash/pwsh: ToolRuntime re-reads the
419 // caller signal right before execute(), and this branch reaches
420 // jobs.start synchronously from there. The shell tools await a
421 // sandbox escalation approval before registering, which is the window
422 // their check covers.
423 return startBackgroundRun(ctx, args, parent, exec.parent === undefined, {
424 recorder,
425 mirror,
426 maxResultChars,
427 })
428 }
429
430 // Meta/body validation failures (META_INVALID/SCRIPT_PARSE) throw
431 // synchronously here and become isError results via the registry — the
432 // model sees the violation list and can correct the call.
433 const run = ctx.workflowEngine.start({
434 script: args.script,
435 meta: args.meta,
436 ...args.args !== undefined ? { args: args.args } : {},
437 parent,
438 signal: exec.signal,
439 })
440 const recordsRun = exec.parent === undefined
441 // The engine publishes member events after start() returns and this run record is active.
442 if (recordsRun) recorder.start(parent.session, run)
443
444 // Bridge the tool's abort signal to the run: if the parent step is aborted while the
445 // script is in flight, cancel the whole run. The signal also enters the engine directly, but
446 // this local bridge preserves the tool contract even if an implementation ignores it.
447 const onAbort = (): void => { run.cancel('parent step aborted') }
448 exec.signal.addEventListener('abort', onAbort, { once: true })
449
450 let result: WorkflowResult | undefined
451 try {
452 result = await run.result
453 const error = stopReasonError(result)
454 if (error !== undefined) {
455 // Map a non-clean finish to an isError result (the registry turns a
456 // throw into an isError). Report the reason, not partial output.
457 throw new Error(error)
458 }
459 return {
460 kind: 'foreground' as const,
461 runId: run.id,
462 agentsStarted: result.agentsStarted,
463 result: result.value as JsonValue,
464 }
465 } finally {
466 exec.signal.removeEventListener('abort', onAbort)
467 try {
468 // Keep member listeners alive through disposal: an engine may
469 // synthesize cancelled member endings while reaching quiescence.
470 await run.dispose()
471 if (recordsRun) {
472 /* v8 ignore next -- WorkflowRun.result never rejects by contract, so result is assigned before finally. */
473 if (result === undefined) throw new Error('workflow run settled without a result')
474 recorder.finish(run.id, result.stopReason)
475 }
476 } finally {
477 if (recordsRun) recorder.abandon(run.id)
478 }
479 }
480 },
481 presentCall: args => presentWorkflowCall(args),
482 presentResult: (args, result) => presentWorkflowResult(args, result),
483 }))
484}