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 propagate4
* through combinators; ordinary child failures and stage errors become per-item nulls.5
* @module @deepseek-ai/dsh-workflow-ptc/runtime6
*/8
import * as vm from 'node:vm'9
import { brandString } from '@deepseek-ai/dsh-brand'10
import type { ContentBlock } from '@deepseek-ai/dsh-llm'11
import type { SessionId } from '@deepseek-ai/dsh-session'12
import { assertObjectJsonSchema, JsonSchemaError } from '@deepseek-ai/dsh-tools'13
import type { ObjectJsonSchema } from '@deepseek-ai/dsh-tools'14
import { isFatalWorkflowError, WorkflowError } from '@deepseek-ai/dsh-workflow'15
import type {16
WorkflowAgentEndInfo,17
WorkflowAgentInfo,18
WorkflowMeta,19
WorkflowResult,20
} from '@deepseek-ai/dsh-workflow'21
import { materializeFromRealm, MaterializeError, renderThrown } from './realm.ts'22
import type { ChildHandle, ChildPort, WorkerLimits } from './types.ts'24
/** The observers the execution reports progress through (the session posts them to the host). */25
export interface ExecutionObserver {26
phase(title: string): void27
log(message: string): void28
agentStart(info: WorkflowAgentInfo): void29
agentEnd(info: WorkflowAgentEndInfo): void30
}32
/** The `agent()` options the script may pass; everything else rejects loud. */33
const SUPPORTED_AGENT_OPTIONS = new Set(['label', 'phase', 'schema', 'provider', 'model'])34
/** Deferred Claude Code options we name explicitly in the rejection message. */35
const DEFERRED_AGENT_OPTIONS = new Set(['effort', 'isolation', 'agentType'])37
/** Flatten a child's final output blocks to text (the non-schema `agent()` result). */38
function outputText(blocks: ContentBlock[]): string {39
return blocks40
.filter((block): block is Extract<ContentBlock, { type: 'text' }> => block.type === 'text')41
.map(block => block.text)42
.join('')43
}45
/** A short display label derived from the prompt when the script passes none. */46
function 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
}52
/**53
* One script execution inside the confined Node process. The host owns54
* cancellation and cleanup of any dropped child work.55
*/56
export class WorkflowExecution {57
/** 1-based count of `agent()` calls started (the `agentsStarted` result field). */58
private started = 059
private activeSlots = 060
private readonly slotWaiters: (() => void)[] = []61
private currentPhase: string | undefined62
private readonly context: vm.Context63
private readonly compiled: vm.Script65
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
}83
this.context = vm.createContext({}, { name: `workflow:${meta.name}` })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) : value98
}99
}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
}116
/**117
* Attach a no-op rejection consumer WITHOUT changing what the caller118
* receives: if the script drops the promise, a host rejection cannot become119
* an unhandled rejection that kills the process; if120
* 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 promise125
}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 error134
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
}142
/** Acquire one concurrency slot in FIFO order. */143
private acquireSlot(): Promise<void> {144
if (this.activeSlots < this.limits.maxConcurrentAgents) {145
this.activeSlots += 1146
return Promise.resolve()147
}148
return new Promise<void>((resolve) => {149
this.slotWaiters.push(() => {150
this.activeSlots += 1151
resolve()152
})153
})154
}156
private releaseSlot(): void {157
this.activeSlots -= 1158
const next = this.slotWaiters.shift()159
if (next) next()160
}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 += 1175
const seq = this.started176
const label = opts.label ?? defaultLabel(rawPrompt)177
const phase = opts.phase ?? this.currentPhase179
await this.acquireSlot()180
try {181
let run: ChildHandle182
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 result196
try {197
result = await run.result198
} catch (error: unknown) {199
// A rejected child result is an INFRASTRUCTURE fault relayed by the200
// host — distinct from a child that failed and resolved. Pair the201
// lifecycle before propagating, and propagate FATAL: an ordinary202
// 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), so210
// 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 null214
}215
this.observer.agentEnd({ ...info, outcome: 'completed' })216
return result.structured217
}218
this.observer.agentEnd({ ...info, outcome: 'completed' })219
return outputText(result.output)220
}221
this.observer.agentEnd({ ...info, outcome: 'failed' })222
return null223
} finally {224
await run.dispose()225
}226
} finally {227
this.releaseSlot()228
}229
}231
/** Materialize + validate the `agent()` options bag from the realm. */232
private readAgentOptions(rawOpts: unknown): {233
label?: string234
phase?: string235
provider?: string236
model?: string237
schema?: ObjectJsonSchema238
} {239
if (rawOpts === undefined) return {}240
let opts: unknown241
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 error246
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)) continue254
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 | undefined265
if (record.schema !== undefined) {266
try {267
assertObjectJsonSchema(record.schema)268
schema = record.schema269
} catch (error: unknown) {270
/* v8 ignore next -- defensive rethrow arm: assertObjectJsonSchema only throws JsonSchemaError */271
if (!(error instanceof JsonSchemaError)) throw error272
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
}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 () => unknown295
})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 be303
// forged (nor accidentally dissolved).304
if (isFatalWorkflowError(error)) throw error305
return null306
}307
}))308
}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) => unknown324
})325
return Promise.all(rawItems.map(async (item: unknown, index) => {326
let value: unknown = item327
try {328
for (const stage of stages) {329
value = await stage(value, item, index)330
}331
return value332
} catch (error: unknown) {333
// An ordinary stage throw drops the ITEM to null and skips its334
// remaining stages; a fatal WorkflowError (see parallel()) kills the335
// whole script.336
if (isFatalWorkflowError(error)) throw error337
return null338
}339
}))340
}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
}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 = title357
this.observer.phase(title)358
}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
}