返回源码地图

packages/workflow/workflow-ptc/src/runtime.ts

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

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

1/**
2 * Workflow VM hooks, child callbacks, ordinary concurrency limits and result serialization.
3 * PTC owns process confinement and cancellation. Fatal hook and provider failures propagate
4 * through combinators; ordinary child failures and stage errors become per-item nulls.
5 * @module @deepseek-ai/dsh-workflow-ptc/runtime
6 */
7
8import * as vm from 'node:vm'
9import { brandString } from '@deepseek-ai/dsh-brand'
10import type { ContentBlock } from '@deepseek-ai/dsh-llm'
11import type { SessionId } from '@deepseek-ai/dsh-session'
12import { assertObjectJsonSchema, JsonSchemaError } from '@deepseek-ai/dsh-tools'
13import type { ObjectJsonSchema } from '@deepseek-ai/dsh-tools'
14import { isFatalWorkflowError, WorkflowError } from '@deepseek-ai/dsh-workflow'
15import type {
16 WorkflowAgentEndInfo,
17 WorkflowAgentInfo,
18 WorkflowMeta,
19 WorkflowResult,
20} from '@deepseek-ai/dsh-workflow'
21import { materializeFromRealm, MaterializeError, renderThrown } from './realm.ts'
22import type { ChildHandle, ChildPort, WorkerLimits } from './types.ts'
23
24/** The observers the execution reports progress through (the session posts them to the host). */
25export interface ExecutionObserver {
26 phase(title: string): void
27 log(message: string): void
28 agentStart(info: WorkflowAgentInfo): void
29 agentEnd(info: WorkflowAgentEndInfo): void
30}
31
32/** The `agent()` options the script may pass; everything else rejects loud. */
33const SUPPORTED_AGENT_OPTIONS = new Set(['label', 'phase', 'schema', 'provider', 'model'])
34/** Deferred Claude Code options we name explicitly in the rejection message. */
35const DEFERRED_AGENT_OPTIONS = new Set(['effort', 'isolation', 'agentType'])
36
37/** Flatten a child's final output blocks to text (the non-schema `agent()` result). */
38function outputText(blocks: ContentBlock[]): string {
39 return blocks
40 .filter((block): block is Extract<ContentBlock, { type: 'text' }> => block.type === 'text')
41 .map(block => block.text)
42 .join('')
43}
44
45/** A short display label derived from the prompt when the script passes none. */
46function defaultLabel(prompt: string): string {
47 const newline = prompt.indexOf('\n')
48 const line = newline === -1 ? prompt : prompt.slice(0, newline)
49 return line.length <= 48 ? line : `${line.slice(0, 47)}…`
50}
51
52/**
53 * One script execution inside the confined Node process. The host owns
54 * cancellation and cleanup of any dropped child work.
55 */
56export class WorkflowExecution {
57 /** 1-based count of `agent()` calls started (the `agentsStarted` result field). */
58 private started = 0
59 private activeSlots = 0
60 private readonly slotWaiters: (() => void)[] = []
61 private currentPhase: string | undefined
62 private readonly context: vm.Context
63 private readonly compiled: vm.Script
64
65 constructor(
66 meta: WorkflowMeta,
67 body: string,
68 args: unknown,
69 private readonly limits: WorkerLimits,
70 private readonly observer: ExecutionObserver,
71 private readonly children: ChildPort,
72 ) {
73 // The host parses this same wrapper before publishing a run.
74 try {
75 this.compiled = new vm.Script(`(async () => {\n${body}\n})()`, {
76 filename: `workflow:${meta.name}`,
77 lineOffset: -1,
78 })
79 } catch (error: unknown) {
80 throw new WorkflowError(`workflow script does not parse: ${String(error)}`, 'SCRIPT_PARSE', { cause: error })
81 }
82
83 this.context = vm.createContext({}, { name: `workflow:${meta.name}` })
84
85 const globals: Record<string, unknown> = {
86 agent: (prompt: unknown, opts?: unknown) => this.contain(this.agent(prompt, opts)),
87 parallel: (thunks: unknown) => this.contain(this.parallel(thunks)),
88 pipeline: (items: unknown, ...stages: unknown[]) => this.contain(this.pipeline(items, stages)),
89 phase: (title: unknown) => { this.phase(title) },
90 log: (message: unknown) => { this.log(message) },
91 // PTC has already copied these inputs through its JSON channel.
92 args,
93 }
94 for (const [key, value] of Object.entries(globals)) {
95 // Data properties on the contextified global; frozen shape not required —
96 // a script overwriting its own hooks only sabotages itself.
97 ;(this.context as Record<string, unknown>)[key] = typeof value === 'function' ? Object.freeze(value) : value
98 }
99 }
100
101 /**
102 * Run the script and materialize its JSON return value.
103 * @returns A completed or error result; script failures never reject.
104 */
105 async drive(): Promise<WorkflowResult> {
106 try {
107 const scriptPromise = this.compiled.runInContext(this.context, { timeout: this.limits.syncTimeoutMs }) as Promise<unknown>
108 const raw: unknown = await this.contain(Promise.resolve(scriptPromise))
109 const value = raw === undefined ? null : this.materializeResult(raw)
110 return { value, stopReason: 'completed', agentsStarted: this.started }
111 } catch (error: unknown) {
112 return { value: null, stopReason: 'error', error: renderThrown(error), agentsStarted: this.started }
113 }
114 }
115
116 /**
117 * Attach a no-op rejection consumer WITHOUT changing what the caller
118 * receives: if the script drops the promise, a host rejection cannot become
119 * an unhandled rejection that kills the process; if
120 * the script does await it, it still observes the rejection.
121 */
122 private contain<T>(promise: Promise<T>): Promise<T> {
123 promise.catch(() => { /* consumed: see method contract — a dropped hook promise must not surface an unhandled rejection */ })
124 return promise
125 }
126
127 /** Materialize the script's return value; violations become RESULT_UNSERIALIZABLE. */
128 private materializeResult(raw: unknown): unknown {
129 try {
130 return materializeFromRealm(raw, 'workflow result')
131 } catch (error: unknown) {
132 /* v8 ignore next -- defensive rethrow arm: materializeFromRealm only throws MaterializeError */
133 if (!(error instanceof MaterializeError)) throw error
134 throw new WorkflowError(
135 `the workflow's return value is not plain JSON data — ${error.message}. Return only JSON-serializable objects/arrays/scalars.`,
136 'RESULT_UNSERIALIZABLE',
137 { cause: error },
138 )
139 }
140 }
141
142 /** Acquire one concurrency slot in FIFO order. */
143 private acquireSlot(): Promise<void> {
144 if (this.activeSlots < this.limits.maxConcurrentAgents) {
145 this.activeSlots += 1
146 return Promise.resolve()
147 }
148 return new Promise<void>((resolve) => {
149 this.slotWaiters.push(() => {
150 this.activeSlots += 1
151 resolve()
152 })
153 })
154 }
155
156 private releaseSlot(): void {
157 this.activeSlots -= 1
158 const next = this.slotWaiters.shift()
159 if (next) next()
160 }
161
162 /** The `agent(prompt, opts)` hook. */
163 private async agent(rawPrompt: unknown, rawOpts: unknown): Promise<unknown> {
164 if (typeof rawPrompt !== 'string' || rawPrompt.length === 0) {
165 throw new WorkflowError('agent() requires a non-empty prompt string', 'INVALID_ARGUMENT')
166 }
167 const opts = this.readAgentOptions(rawOpts)
168 if (this.started >= this.limits.maxTotalAgents) {
169 throw new WorkflowError(
170 `this run reached its total agent cap (${this.limits.maxTotalAgents}) — a runaway-loop backstop; raise the applicable maxTotalAgents limit if the scale is intentional`,
171 'AGENT_CAP',
172 )
173 }
174 this.started += 1
175 const seq = this.started
176 const label = opts.label ?? defaultLabel(rawPrompt)
177 const phase = opts.phase ?? this.currentPhase
178
179 await this.acquireSlot()
180 try {
181 let run: ChildHandle
182 try {
183 run = await this.children.startAgent({
184 prompt: rawPrompt,
185 ...opts.schema !== undefined ? { schema: opts.schema } : {},
186 ...opts.provider !== undefined ? { provider: opts.provider } : {},
187 ...opts.model !== undefined ? { model: opts.model } : {},
188 })
189 } catch (error: unknown) {
190 throw new WorkflowError(`agent() could not start a child: ${renderThrown(error)}`, 'AGENT_START', { cause: error })
191 }
192 const info: WorkflowAgentInfo = { seq, label, ...phase !== undefined ? { phase } : {}, childId: brandString<SessionId>(run.id) }
193 this.observer.agentStart(info)
194 try {
195 let result
196 try {
197 result = await run.result
198 } catch (error: unknown) {
199 // A rejected child result is an INFRASTRUCTURE fault relayed by the
200 // host — distinct from a child that failed and resolved. Pair the
201 // lifecycle before propagating, and propagate FATAL: an ordinary
202 // throw would dissolve to a per-item null inside the combinators,
203 // and a broken provider must not read as a failed child.
204 this.observer.agentEnd({ ...info, outcome: 'failed' })
205 throw new WorkflowError(`child agent run failed: ${renderThrown(error)}`, 'AGENT_RESULT', { cause: error })
206 }
207 if (result.stopReason === 'completed') {
208 if (opts.schema !== undefined) {
209 // The provider honored outputSchema (capability-gated at start), so
210 // a completed run without a structured value is a child failure.
211 if (result.structured === undefined) {
212 this.observer.agentEnd({ ...info, outcome: 'failed' })
213 return null
214 }
215 this.observer.agentEnd({ ...info, outcome: 'completed' })
216 return result.structured
217 }
218 this.observer.agentEnd({ ...info, outcome: 'completed' })
219 return outputText(result.output)
220 }
221 this.observer.agentEnd({ ...info, outcome: 'failed' })
222 return null
223 } finally {
224 await run.dispose()
225 }
226 } finally {
227 this.releaseSlot()
228 }
229 }
230
231 /** Materialize + validate the `agent()` options bag from the realm. */
232 private readAgentOptions(rawOpts: unknown): {
233 label?: string
234 phase?: string
235 provider?: string
236 model?: string
237 schema?: ObjectJsonSchema
238 } {
239 if (rawOpts === undefined) return {}
240 let opts: unknown
241 try {
242 opts = materializeFromRealm(rawOpts, 'agent() options')
243 } catch (error: unknown) {
244 /* v8 ignore next -- defensive rethrow arm: materializeFromRealm only throws MaterializeError */
245 if (!(error instanceof MaterializeError)) throw error
246 throw new WorkflowError(`agent() options must be plain JSON data — ${error.message}`, 'INVALID_ARGUMENT', { cause: error })
247 }
248 if (typeof opts !== 'object' || opts === null || Array.isArray(opts)) {
249 throw new WorkflowError('agent() options must be an object', 'INVALID_ARGUMENT')
250 }
251 const record = opts as Record<string, unknown>
252 for (const key of Object.keys(record)) {
253 if (SUPPORTED_AGENT_OPTIONS.has(key)) continue
254 if (DEFERRED_AGENT_OPTIONS.has(key)) {
255 throw new WorkflowError(`agent() option "${key}" is deferred and not supported by this engine (supported: label, phase, schema, provider, model)`, 'UNSUPPORTED_OPTION')
256 }
257 throw new WorkflowError(`agent() option "${key}" is not recognized (supported: label, phase, schema, provider, model)`, 'UNSUPPORTED_OPTION')
258 }
259 for (const key of ['label', 'phase', 'provider', 'model'] as const) {
260 if (record[key] !== undefined && typeof record[key] !== 'string') {
261 throw new WorkflowError(`agent() option "${key}" must be a string`, 'INVALID_ARGUMENT')
262 }
263 }
264 let schema: ObjectJsonSchema | undefined
265 if (record.schema !== undefined) {
266 try {
267 assertObjectJsonSchema(record.schema)
268 schema = record.schema
269 } catch (error: unknown) {
270 /* v8 ignore next -- defensive rethrow arm: assertObjectJsonSchema only throws JsonSchemaError */
271 if (!(error instanceof JsonSchemaError)) throw error
272 throw new WorkflowError(`agent() schema is outside the supported subset — ${error.message}`, 'UNSUPPORTED_SCHEMA', { cause: error })
273 }
274 }
275 return {
276 ...record.label !== undefined ? { label: record.label as string } : {},
277 ...record.phase !== undefined ? { phase: record.phase as string } : {},
278 ...record.provider !== undefined ? { provider: record.provider as string } : {},
279 ...record.model !== undefined ? { model: record.model as string } : {},
280 ...schema !== undefined ? { schema } : {},
281 }
282 }
283
284 /** The `parallel(thunks)` hook: each thunk caught → `null`; fatal errors propagate. */
285 private async parallel(rawThunks: unknown): Promise<unknown[]> {
286 if (!Array.isArray(rawThunks)) {
287 throw new WorkflowError('parallel() requires an array of zero-argument functions', 'INVALID_ARGUMENT')
288 }
289 this.assertItemCap(rawThunks.length, 'parallel()')
290 const thunks = rawThunks.map((thunk, index) => {
291 if (typeof thunk !== 'function') {
292 throw new WorkflowError(`parallel() item ${index} is not a function`, 'INVALID_ARGUMENT')
293 }
294 return thunk as () => unknown
295 })
296 return Promise.all(thunks.map(async (thunk) => {
297 try {
298 return await thunk()
299 } catch (error: unknown) {
300 // Hook failures are WorkflowErrors built OUTSIDE the script's realm;
301 // fatality is recognized by `instanceof` against this realm's class —
302 // a script-built object can never pass it, so fatality cannot be
303 // forged (nor accidentally dissolved).
304 if (isFatalWorkflowError(error)) throw error
305 return null
306 }
307 }))
308 }
309
310 /** The `pipeline(items, ...stages)` hook: per-item stage chains, NO cross-stage barrier. */
311 private async pipeline(rawItems: unknown, rawStages: unknown[]): Promise<unknown[]> {
312 if (!Array.isArray(rawItems)) {
313 throw new WorkflowError('pipeline() requires an items array', 'INVALID_ARGUMENT')
314 }
315 this.assertItemCap(rawItems.length, 'pipeline()')
316 if (rawStages.length === 0) {
317 throw new WorkflowError('pipeline() requires at least one stage function', 'INVALID_ARGUMENT')
318 }
319 const stages = rawStages.map((stage, index) => {
320 if (typeof stage !== 'function') {
321 throw new WorkflowError(`pipeline() stage ${index} is not a function`, 'INVALID_ARGUMENT')
322 }
323 return stage as (previous: unknown, item: unknown, index: number) => unknown
324 })
325 return Promise.all(rawItems.map(async (item: unknown, index) => {
326 let value: unknown = item
327 try {
328 for (const stage of stages) {
329 value = await stage(value, item, index)
330 }
331 return value
332 } catch (error: unknown) {
333 // An ordinary stage throw drops the ITEM to null and skips its
334 // remaining stages; a fatal WorkflowError (see parallel()) kills the
335 // whole script.
336 if (isFatalWorkflowError(error)) throw error
337 return null
338 }
339 }))
340 }
341
342 private assertItemCap(length: number, hook: string): void {
343 if (length > this.limits.maxItemsPerCall) {
344 throw new WorkflowError(
345 `${hook} received ${length} items — over the per-call cap (${this.limits.maxItemsPerCall}); split the work or raise maxItemsPerCall in the engine config`,
346 'ITEM_CAP',
347 )
348 }
349 }
350
351 /** The `phase(title)` hook: sets the current label for subsequent `agent()` calls and notifies observers. */
352 private phase(title: unknown): void {
353 if (typeof title !== 'string' || title.length === 0) {
354 throw new WorkflowError('phase() requires a non-empty title string', 'INVALID_ARGUMENT')
355 }
356 this.currentPhase = title
357 this.observer.phase(title)
358 }
359
360 /** The `log(message)` hook: narration to observers. */
361 private log(message: unknown): void {
362 if (typeof message !== 'string') {
363 throw new WorkflowError('log() requires a message string', 'INVALID_ARGUMENT')
364 }
365 this.observer.log(message)
366 }
367}