1
/** Workflow child ownership and progress over the shared sandboxed PTC executor. */2
import type { Context } from '@deepseek-ai/cordis'3
import type { Agent } from '@deepseek-ai/dsh-agent'4
import type { PtcBindingFunction, PtcJsonValue, PtcRuntime } from '@deepseek-ai/dsh-ptc-runtime'5
import type { SandboxExecutionPolicy } from '@deepseek-ai/dsh-sandbox'6
import { SessionId } from '@deepseek-ai/dsh-session'7
import type SubagentRuntime from '@deepseek-ai/dsh-subagent'8
import type { SubagentRun } from '@deepseek-ai/dsh-subagent'9
import { assertObjectJsonSchema } from '@deepseek-ai/dsh-tools'10
import type { ObjectJsonSchema } from '@deepseek-ai/dsh-tools'11
import { assertNever, snapshotJsonValue } from '@deepseek-ai/dsh-util-values'12
import type { WorkflowAgentEndInfo, WorkflowAgentInfo, WorkflowMeta, WorkflowResult, WorkflowRun, WorkflowRunId } from '@deepseek-ai/dsh-workflow'13
import { WORKFLOW_GUEST_SOURCE } from './guest-source.ts'14
import type { WorkflowProgress } from './guest-types.ts'15
import { renderThrown } from './realm.ts'16
import type { ExecutionObserver } from './runtime.ts'17
import type { ChildStartRequest, WorkerInit } from './types.ts'19
interface ChildRecord {20
readonly callId: number21
readonly run: SubagentRun22
disposal?: Promise<void>23
}25
const GUEST_URL = `data:text/javascript,${encodeURIComponent(WORKFLOW_GUEST_SOURCE)}`26
const PROGRAM = `const { runWorkflowGuest } = await import(${JSON.stringify(GUEST_URL)}); return await runWorkflowGuest(workflowHost);`28
function 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
}33
function text(value: unknown, name: string): string {34
if (typeof value !== 'string') throw new Error(`workflow ${name} must be a string`)35
return value36
}38
function 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 PtcJsonValue42
}44
function 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 | undefined50
if (request.schema !== undefined) {51
const candidate = object(request.schema)52
assertObjectJsonSchema(candidate)53
schema = candidate54
}55
return {56
prompt,57
...provider === undefined ? {} : { provider },58
...model === undefined ? {} : { model },59
...schema === undefined ? {} : { schema },60
}61
}63
function 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
}74
function 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
}89
function 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
}94
function 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
}107
/**108
* Holder-owned workflow. Cancellation stops the program immediately; settlement waits for109
* its managed process and every admitted child startup/disposal. Engine unload does not110
* invalidate the captured runtime or subagent handles.111
*/112
export 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 = 0119
private terminal = false120
private cancelReason: string | undefined121
private disposed: Promise<void> | undefined122
private readonly externalAbort: () => void124
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
}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) return150
this.cancelReason = reason151
this.controller.abort(reason)152
for (const record of this.children.values()) void this.disposeChild(record)153
}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.disposed163
}165
private requireActive(): void {166
this.controller.signal.throwIfAborted()167
}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 task173
}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
}188
private child(value: unknown): ChildRecord {189
this.requireActive()190
const callId = object(value).callId191
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 record195
}197
private async startChild(request: ChildStartRequest): Promise<PtcJsonValue> {198
this.requireActive()199
const callId = ++this.started200
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
}222
private async childResult(record: ChildRecord): Promise<PtcJsonValue> {223
const signal = this.controller.signal224
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
}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.disposal247
}249
private onProgress(event: WorkflowProgress): void {250
this.requireActive()251
switch (event.type) {252
case 'phase': this.observer.phase(event.title); break253
case 'log': this.observer.log(event.message); break254
case 'agent-start':255
this.liveAgents.set(event.info.seq, event.info)256
this.observer.agentStart(event.info)257
break258
case 'agent-end': this.endAgent(event.info); break259
/* v8 ignore next -- progress() validates the closed message union before dispatch. */260
default: assertNever(event, 'workflow progress')261
}262
}264
private endAgent(info: WorkflowAgentEndInfo): void {265
if (!this.liveAgents.delete(info.seq)) return266
this.observer.agentEnd(info)267
}269
private cancelled(): WorkflowResult {270
return { value: null, stopReason: 'cancelled', error: `workflow run cancelled: ${this.cancelReason}`, agentsStarted: this.started }271
}273
private async drive(): Promise<WorkflowResult> {274
let result: WorkflowResult275
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 = true285
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 = true290
result = this.cancelReason === undefined291
? { value: null, stopReason: 'error', error: renderThrown(error), agentsStarted: this.started }292
: this.cancelled()293
} finally {294
this.terminal = true295
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 result305
}306
}