runtime.ts 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504
  1. /**
  2. * Per-run execution state for the engine's THREAD side: the script's vm
  3. * context and its injected hooks (`agent`/`parallel`/`pipeline`/`phase`/
  4. * `log`/`args`), the concurrency semaphore and caps, cancellation, and the
  5. * drive loop that turns a script settlement into a {@link WorkflowResult}.
  6. * Children are started by RPC to the host through a {@link ChildPort}, so
  7. * this module never touches a cordis context — it runs inside the worker
  8. * thread.
  9. *
  10. * Value boundary (the trust premise lives in ./realm.ts): values ENTERING the
  11. * worker-side host code from the script (hook options, schemas, the return
  12. * value) are materialized by `materializeFromRealm` — a plain walk that
  13. * rejects loud everything JSON cannot carry, which also makes every value
  14. * safe for the later postMessage hop. Values ENTERING the realm (`args`,
  15. * `agent()` results, hook promises and their failures, combinator arrays) are
  16. * handed over DIRECTLY as worker-realm values: the script is model-written
  17. * and trusted, so outer prototypes are not a leak. `args` is cloned once at
  18. * start so a script scribbling on it cannot mutate the session's init object
  19. * (a benign-bug guard; the postMessage clone already isolated the caller).
  20. *
  21. * Failure discipline: fatal {@link WorkflowError}s (bad hook arguments,
  22. * unsupported options/schemas, tripped caps, synchronous start refusal,
  23. * provider-start failure, ready-child result rejection, and
  24. * cancellation) ALWAYS propagate through
  25. * `parallel`/`pipeline` — recognized by `instanceof` against this realm's
  26. * class, which a script inside the vm context cannot forge — and the per-item
  27. * `null` is reserved for child-run failures and ordinary in-stage script
  28. * errors. Every hook-returned promise gets a no-op rejection consumer, so a
  29. * dropped promise cannot surface an unhandled rejection (which would kill the
  30. * worker and read as an engine fault).
  31. *
  32. * There is deliberately NO worker-side abandon channel: a script that never
  33. * settles after a cancel simply never posts a result, and the HOST enforces
  34. * the settles-within-grace guarantee by force-settling `cancelled` and
  35. * terminating the worker — the real kill an in-process engine could not have.
  36. *
  37. * @module @deepseek-ai/dsh-workflow-workerthread/runtime
  38. */
  39. import * as vm from 'node:vm'
  40. import { AgentId } from '@deepseek-ai/dsh-agent'
  41. import type { ContentBlock } from '@deepseek-ai/dsh-llm'
  42. import { assertSupportedOutputSchema, OutputSchemaError } from '@deepseek-ai/dsh-tools'
  43. import type { StructuredOutputSchema } from '@deepseek-ai/dsh-tools'
  44. import { isFatalWorkflowError, WorkflowError } from '@deepseek-ai/dsh-workflow'
  45. import type {
  46. WorkflowAgentEndInfo,
  47. WorkflowAgentInfo,
  48. WorkflowMeta,
  49. WorkflowResult,
  50. } from '@deepseek-ai/dsh-workflow'
  51. import { materializeFromRealm, MaterializeError, renderThrown } from './realm.ts'
  52. import type { ChildHandle, ChildPort, WorkerLimits } from './types.ts'
  53. /** The observers the execution reports progress through (the session posts them to the host). */
  54. export interface ExecutionObserver {
  55. phase(title: string): void
  56. log(message: string): void
  57. agentStart(info: WorkflowAgentInfo): void
  58. agentEnd(info: WorkflowAgentEndInfo): void
  59. }
  60. /** The `agent()` options the script may pass; everything else rejects loud. */
  61. const SUPPORTED_AGENT_OPTIONS = new Set(['label', 'phase', 'schema', 'model'])
  62. /** Deferred Claude Code options we name explicitly in the rejection message. */
  63. const DEFERRED_AGENT_OPTIONS = new Set(['effort', 'isolation', 'agentType'])
  64. /** Flatten a child's final output blocks to text (the non-schema `agent()` result). */
  65. function outputText(blocks: ContentBlock[]): string {
  66. return blocks
  67. .filter((block): block is Extract<ContentBlock, { type: 'text' }> => block.type === 'text')
  68. .map(block => block.text)
  69. .join('')
  70. }
  71. /** A short display label derived from the prompt when the script passes none. */
  72. function defaultLabel(prompt: string): string {
  73. const newline = prompt.indexOf('\n')
  74. const line = newline === -1 ? prompt : prompt.slice(0, newline)
  75. return line.length <= 48 ? line : `${line.slice(0, 47)}…`
  76. }
  77. /**
  78. * One live script execution inside the worker. Constructed per run by the
  79. * session; `drive()` is called exactly once and NEVER rejects — every failure
  80. * becomes a {@link WorkflowResult} with a non-`completed` stop reason. The
  81. * host owns cancellation and cleanup of any dropped child work.
  82. */
  83. export class WorkflowExecution {
  84. /** 1-based count of `agent()` calls started (the `agentsStarted` result field). */
  85. private started = 0
  86. private activeSlots = 0
  87. private readonly slotWaiters: { resolve(): void; reject(error: unknown): void }[] = []
  88. private cancelReason: string | undefined
  89. private cancelError: WorkflowError | undefined
  90. private currentPhase: string | undefined
  91. private readonly context: vm.Context
  92. private readonly compiled: vm.Script
  93. constructor(
  94. meta: WorkflowMeta,
  95. body: string,
  96. args: unknown,
  97. private readonly limits: WorkerLimits,
  98. private readonly observer: ExecutionObserver,
  99. private readonly children: ChildPort,
  100. ) {
  101. // Compile FIRST: a body syntax error must throw out of the constructor
  102. // before any realm state exists. The host pre-parses the identical
  103. // wrapper, so under one Node version this throw is unreachable in
  104. // production — the session still maps it to an error result defensively.
  105. // lineOffset compensates for the wrapper line, so stack traces carry the
  106. // script's own line numbers.
  107. try {
  108. this.compiled = new vm.Script(`(async () => {\n${body}\n})()`, {
  109. filename: `workflow:${meta.name}`,
  110. lineOffset: -1,
  111. })
  112. } catch (error: unknown) {
  113. throw new WorkflowError(`workflow script does not parse: ${String(error)}`, 'SCRIPT_PARSE', { cause: error })
  114. }
  115. this.context = vm.createContext({}, { name: `workflow:${meta.name}` })
  116. const globals: Record<string, unknown> = {
  117. agent: (prompt: unknown, opts?: unknown) => this.contain(this.agent(prompt, opts)),
  118. parallel: (thunks: unknown) => this.contain(this.parallel(thunks)),
  119. pipeline: (items: unknown, ...stages: unknown[]) => this.contain(this.pipeline(items, stages)),
  120. phase: (title: unknown) => { this.phase(title) },
  121. log: (message: unknown) => { this.log(message) },
  122. // workerData already performed the real cross-thread structured clone.
  123. args,
  124. }
  125. for (const [key, value] of Object.entries(globals)) {
  126. // Data properties on the contextified global; frozen shape not required —
  127. // a script overwriting its own hooks only sabotages itself.
  128. ;(this.context as Record<string, unknown>)[key] = typeof value === 'function' ? Object.freeze(value) : value
  129. }
  130. }
  131. /**
  132. * Whether the run has been cancelled. A METHOD, not an inline property
  133. * read: `cancel()` mutates `cancelReason` concurrently (the session's
  134. * message handler), and an inline read after an `await` gets narrowed by
  135. * control flow into an always-false comparison.
  136. */
  137. private isCancelled(): boolean {
  138. return this.cancelReason !== undefined
  139. }
  140. /**
  141. * Shared hook entry guard: after {@link cancel}, EVERY hook throws
  142. * `CANCELLED` at its next call — cancellation is the next HOOK boundary,
  143. * not just the next `agent()`, so a script that caught one cancelled
  144. * rejection cannot keep emitting progress through `phase`/`log` or enter a
  145. * combinator.
  146. */
  147. private throwIfCancelled(): void {
  148. if (this.isCancelled()) throw this.cancelledError()
  149. }
  150. /**
  151. * Cancel the run: waiting `agent()` slots reject and every future hook call
  152. * throws `CANCELLED` — the script dies at its next await. A script that
  153. * never settles anyway (parked on a promise no hook owns) is the HOST's
  154. * problem: its grace timer force-settles the run and terminates the
  155. * worker. Idempotent; the first reason wins.
  156. * @param reason - human-readable cause carried on the CANCELLED error. The
  157. * host independently aborts the required signal shared by every child.
  158. */
  159. cancel(reason: string): void {
  160. if (this.cancelReason !== undefined) return
  161. this.cancelReason = reason
  162. this.cancelError = new WorkflowError(`workflow run cancelled: ${this.cancelReason}`, 'CANCELLED')
  163. for (const waiter of this.slotWaiters.splice(0)) waiter.reject(this.cancelledError())
  164. }
  165. /**
  166. * Run the script to settlement. Resolves — never rejects — with the run's
  167. * {@link WorkflowResult}: the materialized return value on `completed`, the
  168. * failure message on `error`, and `cancelled` when the script died of
  169. * cancellation. This method only chooses the result; the session publishes
  170. * it and the host owns terminal child cancellation.
  171. * @returns the settled outcome — this promise NEVER rejects (the seam's
  172. * `result`-never-rejects contract); every failure maps to a variant.
  173. */
  174. async drive(): Promise<WorkflowResult> {
  175. try {
  176. // Cancelled before the body ever ran (an already-aborted start signal,
  177. // relayed by the host before its `go`): the script must not execute at
  178. // all, let alone report `completed`.
  179. if (this.isCancelled()) throw this.cancelledError()
  180. const scriptPromise = this.compiled.runInContext(this.context, { timeout: this.limits.syncTimeoutMs }) as Promise<unknown>
  181. const raw: unknown = await this.contain(Promise.resolve(scriptPromise))
  182. // Cancelled while the body ran: a script that settled without touching
  183. // another hook (or without any) must still report `cancelled` — the
  184. // holder asked for cancellation and `completed` would be a lie.
  185. if (this.isCancelled()) throw this.cancelledError()
  186. const value = raw === undefined ? null : this.materializeResult(raw)
  187. return { value, stopReason: 'completed', agentsStarted: this.started }
  188. } catch (error: unknown) {
  189. // Any failure after cancel() reports `cancelled` with the canonical
  190. // reason — the reject path mirrors the resolve path's post-settle check.
  191. if (this.isCancelled()) {
  192. return { value: null, stopReason: 'cancelled', error: this.cancelledError().message, agentsStarted: this.started }
  193. }
  194. // renderThrown is total (thrown values of any realm), so this arm
  195. // cannot throw — drive() resolving is the `result` never-rejects seam
  196. // contract.
  197. return { value: null, stopReason: 'error', error: renderThrown(error), agentsStarted: this.started }
  198. }
  199. }
  200. /**
  201. * Attach a no-op rejection consumer WITHOUT changing what the caller
  202. * receives: if the script drops the promise (no await), cancellation cannot
  203. * become an unhandled rejection (which would kill the worker thread); if
  204. * the script does await it, it still observes the rejection.
  205. */
  206. private contain<T>(promise: Promise<T>): Promise<T> {
  207. promise.catch(() => { /* consumed: see method contract — a dropped hook promise must not surface an unhandled rejection */ })
  208. return promise
  209. }
  210. private cancelledError(): WorkflowError {
  211. // cancel() arms cancelError before any caller can observe isCancelled()
  212. // === true; the fallback guards the type, not a reachable path.
  213. /* v8 ignore next */
  214. return this.cancelError ?? new WorkflowError('workflow run cancelled', 'CANCELLED')
  215. }
  216. /** Materialize the script's return value; violations become RESULT_UNSERIALIZABLE. */
  217. private materializeResult(raw: unknown): unknown {
  218. try {
  219. return materializeFromRealm(raw, 'workflow result')
  220. } catch (error: unknown) {
  221. /* v8 ignore next -- defensive rethrow arm: materializeFromRealm only throws MaterializeError */
  222. if (!(error instanceof MaterializeError)) throw error
  223. throw new WorkflowError(
  224. `the workflow's return value is not plain JSON data — ${error.message}. Return only JSON-serializable objects/arrays/scalars.`,
  225. 'RESULT_UNSERIALIZABLE',
  226. { cause: error },
  227. )
  228. }
  229. }
  230. /**
  231. * Acquire one concurrency slot (FIFO). Cancellation rejects QUEUED waiters
  232. * (see {@link cancel}); the callers guard their own entry and post-acquire
  233. * windows, so no cancelled-precheck is duplicated here.
  234. */
  235. private acquireSlot(): Promise<void> {
  236. if (this.activeSlots < this.limits.maxConcurrentAgents) {
  237. this.activeSlots += 1
  238. return Promise.resolve()
  239. }
  240. return new Promise<void>((resolve, reject) => {
  241. this.slotWaiters.push({
  242. resolve: () => {
  243. this.activeSlots += 1
  244. resolve()
  245. },
  246. reject,
  247. })
  248. })
  249. }
  250. private releaseSlot(): void {
  251. this.activeSlots -= 1
  252. const next = this.slotWaiters.shift()
  253. if (next) next.resolve()
  254. }
  255. /** The `agent(prompt, opts)` hook. */
  256. private async agent(rawPrompt: unknown, rawOpts: unknown): Promise<unknown> {
  257. this.throwIfCancelled()
  258. if (typeof rawPrompt !== 'string' || rawPrompt.length === 0) {
  259. throw new WorkflowError('agent() requires a non-empty prompt string', 'INVALID_ARGUMENT')
  260. }
  261. const opts = this.readAgentOptions(rawOpts)
  262. if (this.started >= this.limits.maxTotalAgents) {
  263. throw new WorkflowError(
  264. `this run reached its total agent cap (${this.limits.maxTotalAgents}) — a runaway-loop backstop; raise maxTotalAgents in the engine config if the scale is intentional`,
  265. 'AGENT_CAP',
  266. )
  267. }
  268. this.started += 1
  269. const seq = this.started
  270. const label = opts.label ?? defaultLabel(rawPrompt)
  271. const phase = opts.phase ?? this.currentPhase
  272. await this.acquireSlot()
  273. try {
  274. // Re-check after the acquire: the await yields at least one microtask
  275. // tick even when a slot is free, and a queued waiter resumes a tick
  276. // after its release — a cancel() landing in either window must not
  277. // reach the host (which would refuse anyway, but the refusal reads as
  278. // a start failure rather than the cancellation it is).
  279. this.throwIfCancelled()
  280. let run: ChildHandle
  281. try {
  282. run = await this.children.startAgent({
  283. prompt: rawPrompt,
  284. ...opts.schema !== undefined ? { schema: opts.schema } : {},
  285. ...opts.model !== undefined ? { model: opts.model } : {},
  286. })
  287. } catch (error: unknown) {
  288. // The host refuses starts once the run is cancelled — a refusal that
  289. // races our own cancel state must read as the cancellation it is,
  290. // not as a broken seam.
  291. if (this.isCancelled()) throw this.cancelledError()
  292. throw new WorkflowError(`agent() could not start a child: ${renderThrown(error)}`, 'AGENT_START', { cause: error })
  293. }
  294. // The start round-trip yields to the event loop, so a cancel CAN land
  295. // between the host starting the child and this continuation running —
  296. // wind the fresh child down instead of leaving it live behind a dead
  297. // script.
  298. if (this.isCancelled()) {
  299. await run.dispose()
  300. throw this.cancelledError()
  301. }
  302. const info: WorkflowAgentInfo = { seq, label, ...phase !== undefined ? { phase } : {}, childId: AgentId(run.id) }
  303. this.observer.agentStart(info)
  304. try {
  305. let result
  306. try {
  307. result = await run.result
  308. } catch (error: unknown) {
  309. // A rejected child result is an INFRASTRUCTURE fault relayed by the
  310. // host — distinct from a child that failed and resolved. Pair the
  311. // lifecycle before propagating, and propagate FATAL: an ordinary
  312. // throw would dissolve to a per-item null inside the combinators,
  313. // and a broken provider must not read as a failed child.
  314. if (this.isCancelled()) {
  315. this.observer.agentEnd({ ...info, outcome: 'cancelled' })
  316. throw this.cancelledError()
  317. }
  318. this.observer.agentEnd({ ...info, outcome: 'failed' })
  319. throw new WorkflowError(`child agent run failed: ${renderThrown(error)}`, 'AGENT_RESULT', { cause: error })
  320. }
  321. if (result.stopReason === 'completed') {
  322. if (opts.schema !== undefined) {
  323. // The provider honored outputSchema (capability-gated at start), so
  324. // a completed run without a structured value is a child failure.
  325. if (result.structured === undefined) {
  326. this.observer.agentEnd({ ...info, outcome: 'failed' })
  327. return null
  328. }
  329. this.observer.agentEnd({ ...info, outcome: 'completed' })
  330. return result.structured
  331. }
  332. this.observer.agentEnd({ ...info, outcome: 'completed' })
  333. return outputText(result.output)
  334. }
  335. // A cancelled RUN kills the script; a child that failed for its own
  336. // reasons resolves null (scripts .filter(Boolean) per the CC contract).
  337. if (this.isCancelled()) {
  338. this.observer.agentEnd({ ...info, outcome: 'cancelled' })
  339. throw this.cancelledError()
  340. }
  341. this.observer.agentEnd({ ...info, outcome: 'failed' })
  342. return null
  343. } finally {
  344. await run.dispose()
  345. }
  346. } finally {
  347. this.releaseSlot()
  348. }
  349. }
  350. /** Materialize + validate the `agent()` options bag from the realm. */
  351. private readAgentOptions(rawOpts: unknown): { label?: string; phase?: string; model?: string; schema?: StructuredOutputSchema } {
  352. if (rawOpts === undefined) return {}
  353. let opts: unknown
  354. try {
  355. opts = materializeFromRealm(rawOpts, 'agent() options')
  356. } catch (error: unknown) {
  357. /* v8 ignore next -- defensive rethrow arm: materializeFromRealm only throws MaterializeError */
  358. if (!(error instanceof MaterializeError)) throw error
  359. throw new WorkflowError(`agent() options must be plain JSON data — ${error.message}`, 'INVALID_ARGUMENT', { cause: error })
  360. }
  361. if (typeof opts !== 'object' || opts === null || Array.isArray(opts)) {
  362. throw new WorkflowError('agent() options must be an object', 'INVALID_ARGUMENT')
  363. }
  364. const record = opts as Record<string, unknown>
  365. for (const key of Object.keys(record)) {
  366. if (SUPPORTED_AGENT_OPTIONS.has(key)) continue
  367. if (DEFERRED_AGENT_OPTIONS.has(key)) {
  368. throw new WorkflowError(`agent() option "${key}" is deferred and not supported by this engine (supported: label, phase, schema, model)`, 'UNSUPPORTED_OPTION')
  369. }
  370. throw new WorkflowError(`agent() option "${key}" is not recognized (supported: label, phase, schema, model)`, 'UNSUPPORTED_OPTION')
  371. }
  372. for (const key of ['label', 'phase', 'model'] as const) {
  373. if (record[key] !== undefined && typeof record[key] !== 'string') {
  374. throw new WorkflowError(`agent() option "${key}" must be a string`, 'INVALID_ARGUMENT')
  375. }
  376. }
  377. let schema: StructuredOutputSchema | undefined
  378. if (record.schema !== undefined) {
  379. try {
  380. assertSupportedOutputSchema(record.schema)
  381. schema = record.schema
  382. } catch (error: unknown) {
  383. /* v8 ignore next -- defensive rethrow arm: assertSupportedOutputSchema only throws OutputSchemaError */
  384. if (!(error instanceof OutputSchemaError)) throw error
  385. throw new WorkflowError(`agent() schema is outside the supported subset — ${error.message}`, 'UNSUPPORTED_SCHEMA', { cause: error })
  386. }
  387. }
  388. return {
  389. ...record.label !== undefined ? { label: record.label as string } : {},
  390. ...record.phase !== undefined ? { phase: record.phase as string } : {},
  391. ...record.model !== undefined ? { model: record.model as string } : {},
  392. ...schema !== undefined ? { schema } : {},
  393. }
  394. }
  395. /** The `parallel(thunks)` hook: each thunk caught → `null`; fatal errors propagate. */
  396. private async parallel(rawThunks: unknown): Promise<unknown[]> {
  397. this.throwIfCancelled()
  398. if (!Array.isArray(rawThunks)) {
  399. throw new WorkflowError('parallel() requires an array of zero-argument functions', 'INVALID_ARGUMENT')
  400. }
  401. this.assertItemCap(rawThunks.length, 'parallel()')
  402. const thunks = rawThunks.map((thunk, index) => {
  403. if (typeof thunk !== 'function') {
  404. throw new WorkflowError(`parallel() item ${index} is not a function`, 'INVALID_ARGUMENT')
  405. }
  406. return thunk as () => unknown
  407. })
  408. return Promise.all(thunks.map(async (thunk) => {
  409. try {
  410. return await thunk()
  411. } catch (error: unknown) {
  412. // Hook failures are WorkflowErrors built OUTSIDE the script's realm;
  413. // fatality is recognized by `instanceof` against this realm's class —
  414. // a script-built object can never pass it, so fatality cannot be
  415. // forged (nor accidentally dissolved).
  416. if (isFatalWorkflowError(error)) throw error
  417. return null
  418. }
  419. }))
  420. }
  421. /** The `pipeline(items, ...stages)` hook: per-item stage chains, NO cross-stage barrier. */
  422. private async pipeline(rawItems: unknown, rawStages: unknown[]): Promise<unknown[]> {
  423. this.throwIfCancelled()
  424. if (!Array.isArray(rawItems)) {
  425. throw new WorkflowError('pipeline() requires an items array', 'INVALID_ARGUMENT')
  426. }
  427. this.assertItemCap(rawItems.length, 'pipeline()')
  428. if (rawStages.length === 0) {
  429. throw new WorkflowError('pipeline() requires at least one stage function', 'INVALID_ARGUMENT')
  430. }
  431. const stages = rawStages.map((stage, index) => {
  432. if (typeof stage !== 'function') {
  433. throw new WorkflowError(`pipeline() stage ${index} is not a function`, 'INVALID_ARGUMENT')
  434. }
  435. return stage as (previous: unknown, item: unknown, index: number) => unknown
  436. })
  437. return Promise.all(rawItems.map(async (item: unknown, index) => {
  438. let value: unknown = item
  439. try {
  440. for (const stage of stages) {
  441. value = await stage(value, item, index)
  442. }
  443. return value
  444. } catch (error: unknown) {
  445. // An ordinary stage throw drops the ITEM to null and skips its
  446. // remaining stages; a fatal WorkflowError (see parallel()) kills the
  447. // whole script.
  448. if (isFatalWorkflowError(error)) throw error
  449. return null
  450. }
  451. }))
  452. }
  453. private assertItemCap(length: number, hook: string): void {
  454. if (length > this.limits.maxItemsPerCall) {
  455. throw new WorkflowError(
  456. `${hook} received ${length} items — over the per-call cap (${this.limits.maxItemsPerCall}); split the work or raise maxItemsPerCall in the engine config`,
  457. 'ITEM_CAP',
  458. )
  459. }
  460. }
  461. /** The `phase(title)` hook: sets the current label for subsequent `agent()` calls and notifies observers. */
  462. private phase(title: unknown): void {
  463. this.throwIfCancelled()
  464. if (typeof title !== 'string' || title.length === 0) {
  465. throw new WorkflowError('phase() requires a non-empty title string', 'INVALID_ARGUMENT')
  466. }
  467. this.currentPhase = title
  468. this.observer.phase(title)
  469. }
  470. /** The `log(message)` hook: narration to observers. */
  471. private log(message: unknown): void {
  472. this.throwIfCancelled()
  473. if (typeof message !== 'string') {
  474. throw new WorkflowError('log() requires a message string', 'INVALID_ARGUMENT')
  475. }
  476. this.observer.log(message)
  477. }
  478. }