host.ts 13 KB

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