返回源码地图

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

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

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

1/** Workflow child ownership and progress over the shared sandboxed PTC executor. */
2import type { Context } from '@deepseek-ai/cordis'
3import type { Agent } from '@deepseek-ai/dsh-agent'
4import type { PtcBindingFunction, PtcJsonValue, PtcRuntime } from '@deepseek-ai/dsh-ptc-runtime'
5import type { SandboxExecutionPolicy } from '@deepseek-ai/dsh-sandbox'
6import { SessionId } from '@deepseek-ai/dsh-session'
7import type SubagentRuntime from '@deepseek-ai/dsh-subagent'
8import type { SubagentRun } from '@deepseek-ai/dsh-subagent'
9import { assertObjectJsonSchema } from '@deepseek-ai/dsh-tools'
10import type { ObjectJsonSchema } from '@deepseek-ai/dsh-tools'
11import { assertNever, snapshotJsonValue } from '@deepseek-ai/dsh-util-values'
12import type { WorkflowAgentEndInfo, WorkflowAgentInfo, WorkflowMeta, WorkflowResult, WorkflowRun, WorkflowRunId } from '@deepseek-ai/dsh-workflow'
13import { WORKFLOW_GUEST_SOURCE } from './guest-source.ts'
14import type { WorkflowProgress } from './guest-types.ts'
15import { renderThrown } from './realm.ts'
16import type { ExecutionObserver } from './runtime.ts'
17import type { ChildStartRequest, WorkerInit } from './types.ts'
18
19interface ChildRecord {
20 readonly callId: number
21 readonly run: SubagentRun
22 disposal?: Promise<void>
23}
24
25const GUEST_URL = `data:text/javascript,${encodeURIComponent(WORKFLOW_GUEST_SOURCE)}`
26const PROGRAM = `const { runWorkflowGuest } = await import(${JSON.stringify(GUEST_URL)}); return await runWorkflowGuest(workflowHost);`
27
28function object(value: unknown): Record<string, unknown> {
29 if (value === null || typeof value !== 'object' || Array.isArray(value)) throw new Error('workflow binding requires an object')
30 return value as Record<string, unknown>
31}
32
33function text(value: unknown, name: string): string {
34 if (typeof value !== 'string') throw new Error(`workflow ${name} must be a string`)
35 return value
36}
37
38function json(value: unknown): PtcJsonValue {
39 const result = snapshotJsonValue(value)
40 if (result === undefined) throw new Error('workflow binding value must be lossless JSON')
41 return result as PtcJsonValue
42}
43
44function childRequest(value: unknown): ChildStartRequest {
45 const request = object(value)
46 const prompt = text(request.prompt, 'prompt')
47 const provider = request.provider === undefined ? undefined : text(request.provider, 'provider')
48 const model = request.model === undefined ? undefined : text(request.model, 'model')
49 let schema: ObjectJsonSchema | undefined
50 if (request.schema !== undefined) {
51 const candidate = object(request.schema)
52 assertObjectJsonSchema(candidate)
53 schema = candidate
54 }
55 return {
56 prompt,
57 ...provider === undefined ? {} : { provider },
58 ...model === undefined ? {} : { model },
59 ...schema === undefined ? {} : { schema },
60 }
61}
62
63function agentInfo(value: unknown): WorkflowAgentInfo {
64 const info = object(value)
65 if (!Number.isSafeInteger(info.seq) || (info.seq as number) < 1) throw new Error('workflow agent sequence must be a positive integer')
66 return {
67 seq: info.seq as number,
68 label: text(info.label, 'agent label'),
69 childId: SessionId(text(info.childId, 'child id')),
70 ...info.phase === undefined ? {} : { phase: text(info.phase, 'agent phase') },
71 }
72}
73
74function progress(value: unknown): WorkflowProgress {
75 const event = object(value)
76 switch (event.type) {
77 case 'phase': return { type: 'phase', title: text(event.title, 'phase') }
78 case 'log': return { type: 'log', message: text(event.message, 'log') }
79 case 'agent-start': return { type: 'agent-start', info: agentInfo(event.info) }
80 case 'agent-end': {
81 const info = object(event.info)
82 if (info.outcome !== 'completed' && info.outcome !== 'failed' && info.outcome !== 'cancelled') throw new Error('invalid workflow agent outcome')
83 return { type: 'agent-end', info: { ...agentInfo(info), outcome: info.outcome } }
84 }
85 default: throw new Error('invalid workflow progress event')
86 }
87}
88
89function progressBatch(value: unknown): WorkflowProgress[] {
90 if (!Array.isArray(value)) throw new Error('workflow progress requires an array of events')
91 return value.map(progress)
92}
93
94function workflowResult(value: unknown): WorkflowResult {
95 const result = object(value)
96 if (result.stopReason !== 'completed' && result.stopReason !== 'error' && result.stopReason !== 'cancelled') throw new Error('invalid workflow stop reason')
97 if (!Number.isSafeInteger(result.agentsStarted) || (result.agentsStarted as number) < 0) throw new Error('invalid workflow agent count')
98 if (!Object.hasOwn(result, 'value')) throw new Error('workflow result is missing its value')
99 return {
100 value: result.value,
101 stopReason: result.stopReason,
102 agentsStarted: result.agentsStarted as number,
103 ...result.error === undefined ? {} : { error: text(result.error, 'error') },
104 }
105}
106
107/**
108 * Holder-owned workflow. Cancellation stops the program immediately; settlement waits for
109 * its managed process and every admitted child startup/disposal. Engine unload does not
110 * invalidate the captured runtime or subagent handles.
111 */
112export class PtcWorkflowRun implements WorkflowRun {
113 readonly result: Promise<WorkflowResult>
114 private readonly controller = new AbortController()
115 private readonly children = new Map<number, ChildRecord>()
116 private readonly pending = new Set<Promise<unknown>>()
117 private readonly liveAgents = new Map<number, WorkflowAgentInfo>()
118 private started = 0
119 private terminal = false
120 private cancelReason: string | undefined
121 private disposed: Promise<void> | undefined
122 private readonly externalAbort: () => void
123
124 constructor(
125 private readonly ctx: Context,
126 private readonly subagents: SubagentRuntime,
127 private readonly runtime: PtcRuntime,
128 readonly id: WorkflowRunId,
129 readonly meta: WorkflowMeta,
130 private readonly parent: Agent,
131 private readonly init: WorkerInit,
132 private readonly provider: string,
133 private readonly policy: SandboxExecutionPolicy,
134 private readonly observer: ExecutionObserver,
135 private readonly signal?: AbortSignal,
136 ) {
137 this.externalAbort = () => { this.cancel('workflow signal aborted') }
138 if (signal?.aborted) this.externalAbort()
139 else signal?.addEventListener('abort', this.externalAbort, { once: true })
140 // Consumers attach durable run recording after start() returns.
141 this.result = Promise.resolve().then(() => this.drive())
142 }
143
144 /**
145 * Stop the script and abort pending and published children.
146 * @param reason - Human-readable cancellation cause; the first request wins.
147 */
148 cancel(reason = 'workflow cancelled'): void {
149 if (this.terminal || this.cancelReason !== undefined) return
150 this.cancelReason = reason
151 this.controller.abort(reason)
152 for (const record of this.children.values()) void this.disposeChild(record)
153 }
154
155 /**
156 * Cancel unfinished work and await the program and child cleanup.
157 * @returns One shared completion promise for repeated disposal calls.
158 */
159 dispose(): Promise<void> {
160 this.cancel('workflow disposed')
161 this.disposed ??= this.result.then(() => {})
162 return this.disposed
163 }
164
165 private requireActive(): void {
166 this.controller.signal.throwIfAborted()
167 }
168
169 private track<T>(task: Promise<T>): Promise<T> {
170 this.pending.add(task)
171 void task.then(() => { this.pending.delete(task) }, () => { this.pending.delete(task) })
172 return task
173 }
174
175 private bindings(): Record<string, PtcBindingFunction> {
176 return {
177 begin: () => { this.requireActive(); return Promise.resolve(json(this.init)) },
178 startChild: value => this.track(this.startChild(childRequest(value))),
179 childResult: value => this.track(this.childResult(this.child(value))),
180 disposeChild: async (value) => { await this.disposeChild(this.child(value)); return null },
181 progress: (value) => {
182 for (const event of progressBatch(value)) this.onProgress(event)
183 return Promise.resolve(null)
184 },
185 }
186 }
187
188 private child(value: unknown): ChildRecord {
189 this.requireActive()
190 const callId = object(value).callId
191 if (!Number.isSafeInteger(callId)) throw new Error('workflow child call id must be an integer')
192 const record = this.children.get(callId as number)
193 if (record === undefined) throw new Error('workflow child call is not active')
194 return record
195 }
196
197 private async startChild(request: ChildStartRequest): Promise<PtcJsonValue> {
198 this.requireActive()
199 const callId = ++this.started
200 const run = await this.subagents.start(this.provider, {
201 prompt: [{ type: 'text', text: request.prompt }],
202 parent: this.parent,
203 signal: this.controller.signal,
204 ...request.schema === undefined ? {} : { outputSchema: request.schema },
205 ...request.provider === undefined && request.model === undefined ? {} : {
206 agentOptions: {
207 ...request.provider === undefined ? {} : { provider: request.provider },
208 ...request.model === undefined ? {} : { model: request.model },
209 },
210 },
211 })
212 const record: ChildRecord = { callId, run }
213 this.children.set(callId, record)
214 // A provider can publish after the signal fired while startup was pending.
215 if (this.controller.signal.aborted) {
216 await this.disposeChild(record)
217 throw new Error('workflow child started after cancellation')
218 }
219 return { callId, childId: run.id }
220 }
221
222 private async childResult(record: ChildRecord): Promise<PtcJsonValue> {
223 const signal = this.controller.signal
224 signal.throwIfAborted()
225 const aborted = Promise.withResolvers<never>()
226 const onAbort = (): void => { aborted.reject(signal.reason) }
227 signal.addEventListener('abort', onAbort, { once: true })
228 try {
229 const result = await Promise.race([record.run.result, aborted.promise])
230 return json({
231 output: result.output,
232 stopReason: result.stopReason,
233 ...result.structured === undefined ? {} : { structured: result.structured },
234 })
235 } finally {
236 signal.removeEventListener('abort', onAbort)
237 }
238 }
239
240 private disposeChild(record: ChildRecord): Promise<void> {
241 record.disposal ??= Promise.resolve().then(() => record.run.dispose()).catch((error: unknown) => {
242 this.ctx.logger.warn(`workflow-ptc: child dispose failed: ${renderThrown(error)}`)
243 }).finally(() => {
244 this.children.delete(record.callId)
245 })
246 return record.disposal
247 }
248
249 private onProgress(event: WorkflowProgress): void {
250 this.requireActive()
251 switch (event.type) {
252 case 'phase': this.observer.phase(event.title); break
253 case 'log': this.observer.log(event.message); break
254 case 'agent-start':
255 this.liveAgents.set(event.info.seq, event.info)
256 this.observer.agentStart(event.info)
257 break
258 case 'agent-end': this.endAgent(event.info); break
259 /* v8 ignore next -- progress() validates the closed message union before dispatch. */
260 default: assertNever(event, 'workflow progress')
261 }
262 }
263
264 private endAgent(info: WorkflowAgentEndInfo): void {
265 if (!this.liveAgents.delete(info.seq)) return
266 this.observer.agentEnd(info)
267 }
268
269 private cancelled(): WorkflowResult {
270 return { value: null, stopReason: 'cancelled', error: `workflow run cancelled: ${this.cancelReason}`, agentsStarted: this.started }
271 }
272
273 private async drive(): Promise<WorkflowResult> {
274 let result: WorkflowResult
275 try {
276 const outcome = await this.runtime.run(this.runtime.resolve({
277 program: PROGRAM,
278 bindings: [{ global: 'workflowHost', functions: this.bindings() }],
279 cwd: this.policy.workspaceRoot,
280 sandboxPolicy: this.policy,
281 timeoutMs: null,
282 signal: this.controller.signal,
283 }))
284 this.terminal = true
285 if (this.cancelReason !== undefined) result = this.cancelled()
286 else if (outcome.error !== undefined) result = { value: null, stopReason: 'error', error: `workflow execution failed (${outcome.error.kind}): ${outcome.error.message}`, agentsStarted: this.started }
287 else result = workflowResult(outcome.value)
288 } catch (error: unknown) {
289 this.terminal = true
290 result = this.cancelReason === undefined
291 ? { value: null, stopReason: 'error', error: renderThrown(error), agentsStarted: this.started }
292 : this.cancelled()
293 } finally {
294 this.terminal = true
295 this.signal?.removeEventListener('abort', this.externalAbort)
296 this.controller.abort('workflow settled')
297 // Disposing published children releases binding waits; pending starts may publish more.
298 for (const record of this.children.values()) void this.disposeChild(record)
299 while (this.pending.size > 0) await Promise.allSettled([...this.pending])
300 await Promise.all([...this.children.values()].map(record => this.disposeChild(record)))
301 this.children.clear()
302 for (const info of this.liveAgents.values()) this.endAgent({ ...info, outcome: 'cancelled' })
303 }
304 return result
305 }
306}