返回源码地图

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

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

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

1/**
2 * Workflow orchestration through the shared sandboxed Node PTC executor.
3 * The VM supplies script helpers; the process applies the calling Session's file policy.
4 * @module @deepseek-ai/dsh-workflow-ptc
5 */
6
7import { randomUUID } from 'node:crypto'
8import { availableParallelism } from 'node:os'
9import * as vm from 'node:vm'
10import type { Context } from '@deepseek-ai/cordis'
11import type {} from '@deepseek-ai/dsh-ptc-runtime'
12import type {} from '@deepseek-ai/dsh-sandbox-policy'
13import z from '@deepseek-ai/schemastery'
14import WorkflowEngine, { WorkflowError, WorkflowRunId } from '@deepseek-ai/dsh-workflow'
15import type { WorkflowRun, WorkflowRunInfo, WorkflowStartRequest } from '@deepseek-ai/dsh-workflow'
16import { PtcWorkflowRun } from './host.ts'
17import { validateMeta } from './meta.ts'
18import type { WorkerInit, WorkerLimits } from './types.ts'
19
20export { validateMeta } from './meta.ts'
21export { materializeFromRealm, MaterializeError } from './realm.ts'
22export type {
23 ChildHandle,
24 ChildPort,
25 ChildResult,
26 ChildStartRequest,
27 WorkerInit,
28 WorkerLimits,
29} from './types.ts'
30
31/** Plugin config (all optional — `static Config` supplies the defaults). */
32export interface Config {
33 /** The `ctx.subagents` provider children run on (default `spawn`). */
34 provider?: string
35 /** Concurrent `agent()` ceiling; `0` (the default) auto-resolves to `min(16, max(1, cores - 2))`. */
36 maxConcurrentAgents?: number
37 /** Total `agent()` calls one run may start — the runaway-loop backstop (default 1000). */
38 maxTotalAgents?: number
39 /** Items accepted by a single `parallel()`/`pipeline()` call (default 4096). */
40 maxItemsPerCall?: number
41 /** VM timeout for the script's initial synchronous slice (default 5000 ms). */
42 syncTimeoutMs?: number
43}
44
45type ResolvedConfig = Required<Config>
46
47/** A body that still carries the Claude Code-style meta header (meta rides the seam as data here). */
48const META_STATEMENT = /^\s*export\s+const\s+meta\b/
49
50/**
51 * Reject invalid JavaScript synchronously before publishing a workflow run.
52 * The guest compiles the same async wrapper in its own process.
53 */
54function assertBodyParses(body: string, name: string): void {
55 if (META_STATEMENT.test(body)) {
56 throw new WorkflowError('workflow meta rides the `meta` request field, not the script: remove the `export const meta = {...}` statement from the body', 'SCRIPT_PARSE')
57 }
58 try {
59 // Parse only — the script object is discarded, nothing executes.
60 void new vm.Script(`(async () => {\n${body}\n})()`, { filename: `workflow:${name}`, lineOffset: -1 })
61 } catch (error: unknown) {
62 throw new WorkflowError(`workflow script does not parse: ${String(error)}`, 'SCRIPT_PARSE', { cause: error })
63 }
64}
65
66/** Resolve one run's provider route before publishing work. */
67function resolveSubagentProvider(ctx: Context, configured: string, override: string | undefined): string {
68 const provider = override ?? configured
69 if (provider.length === 0 || provider !== provider.trim()) {
70 throw new WorkflowError(
71 'workflow subagentProvider must be a non-empty normalized string',
72 'INVALID_ARGUMENT',
73 )
74 }
75 if (ctx.subagents.getProvider(provider) === undefined) {
76 throw new WorkflowError(`no subagent provider registered for "${provider}"`, 'AGENT_START')
77 }
78 return provider
79}
80
81/** Resolve one run's total-child cap against the engine deployment ceiling. */
82function resolveMaxTotalAgents(requested: number | undefined, ceiling: number): number {
83 if (requested === undefined) return ceiling
84 if (!Number.isSafeInteger(requested) || requested < 1) {
85 throw new WorkflowError('workflow maxTotalAgents must be a positive safe integer', 'INVALID_ARGUMENT')
86 }
87 if (requested > ceiling) {
88 throw new WorkflowError(
89 `workflow maxTotalAgents ${requested} exceeds the engine ceiling ${ceiling}`,
90 'INVALID_ARGUMENT',
91 )
92 }
93 return requested
94}
95
96/**
97 * The PTC-backed workflow engine. `start()` validates the script up front
98 * (meta + a host-side body parse) and returns a {@link WorkflowRun} whose
99 * `result` never rejects; the `workflow/*` events fire around the run per
100 * the seam contract.
101 */
102class PtcWorkflowEngine extends WorkflowEngine {
103 static inject = ['subagents', 'ptcRuntime', 'sandboxPolicy']
104
105 static Config: z<Config> = z.object({
106 provider: z.string().default('spawn'),
107 maxConcurrentAgents: z.natural().default(0),
108 maxTotalAgents: z.natural().min(1).default(1000),
109 maxItemsPerCall: z.natural().min(1).default(4096),
110 syncTimeoutMs: z.natural().min(1).default(5000),
111 })
112
113 private readonly config: ResolvedConfig
114
115 constructor(ctx: Context, config: Config) {
116 super(ctx)
117 if (ctx.ptcRuntime.language !== 'typescript') throw new Error('workflow-ptc requires the Node TypeScript PTC runtime')
118 // schemastery (static Config) has already filled the defaulted fields;
119 // the assertion records that resolution, not a hidden fallback.
120 this.config = config as ResolvedConfig
121 }
122
123 /**
124 * Validate and execute a workflow script in a sandboxed Node process. Throws
125 * {@link WorkflowError} synchronously (`META_INVALID` for a malformed meta
126 * block, `SCRIPT_PARSE` for a body that does not compile) for a request
127 * that cannot begin; once a run is returned, every failure resolves through
128 * `result.stopReason` instead.
129 * @param request - the script body, its meta data and `args`, the parent
130 * agent, and an optional cancel signal.
131 * @returns the live run (its `result` resolves when the script settles).
132 */
133 start(request: WorkflowStartRequest): WorkflowRun {
134 const meta = validateMeta(request.meta)
135 assertBodyParses(request.script, meta.name)
136 const subagentProvider = resolveSubagentProvider(this.ctx, this.config.provider, request.subagentProvider)
137 const maxTotalAgents = resolveMaxTotalAgents(request.maxTotalAgents, this.config.maxTotalAgents)
138 const id = WorkflowRunId(randomUUID())
139 const info: WorkflowRunInfo = { id, meta }
140 const limits: WorkerLimits = {
141 maxConcurrentAgents: this.config.maxConcurrentAgents === 0
142 ? Math.min(16, Math.max(1, availableParallelism() - 2))
143 : this.config.maxConcurrentAgents,
144 maxTotalAgents,
145 maxItemsPerCall: this.config.maxItemsPerCall,
146 syncTimeoutMs: this.config.syncTimeoutMs,
147 }
148 const init: WorkerInit = {
149 meta,
150 body: request.script,
151 ...request.args !== undefined ? { args: structuredClone(request.args) } : {},
152 limits,
153 }
154 // Captured service handles keep a holder-owned run usable after engine unload.
155 const runCtx = this.ctx
156 const subagents = runCtx.subagents
157 const run = new PtcWorkflowRun(
158 runCtx,
159 subagents,
160 runCtx.ptcRuntime,
161 id,
162 meta,
163 request.parent,
164 init,
165 subagentProvider,
166 runCtx.sandboxPolicy.resolve({ session: request.parent.session }),
167 {
168 phase: (title) => { this.emitWorkflowEvent('workflow/phase', info, title) },
169 log: (message) => { this.emitWorkflowEvent('workflow/log', info, message) },
170 agentStart: (agent) => { this.emitWorkflowEvent('workflow/agent-start', info, agent) },
171 agentEnd: (agent) => { this.emitWorkflowEvent('workflow/agent-end', info, agent) },
172 },
173 request.signal,
174 )
175
176 this.emitWorkflowEvent('workflow/start', info)
177 // `workflow/end` fires as the (never-rejecting) result settles, with the
178 // outcome DATA only — the value stays with the run's holder.
179 void run.result.then((settled) => {
180 this.emitWorkflowEvent('workflow/end', info, {
181 stopReason: settled.stopReason,
182 ...settled.error !== undefined ? { error: settled.error } : {},
183 agentsStarted: settled.agentsStarted,
184 })
185 })
186
187 return run
188 }
189}
190
191export default PtcWorkflowEngine