index.ts 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205
  1. /**
  2. * Service Definition for the workflow capability seam. Service providers execute orchestration scripts;
  3. * observe-only lifecycle events never expose run control.
  4. * @module @deepseek-ai/dsh-workflow
  5. */
  6. import { Context, Service } from 'cordis'
  7. import { HarnessError } from '@deepseek-ai/dsh-llm'
  8. import type {
  9. WorkflowAgentEndInfo,
  10. WorkflowAgentInfo,
  11. WorkflowResultInfo,
  12. WorkflowRun,
  13. WorkflowRunInfo,
  14. WorkflowStartRequest,
  15. } from './types.ts'
  16. export { WorkflowRunId } from './types.ts'
  17. export type {
  18. WorkflowAgentEndInfo,
  19. WorkflowAgentInfo,
  20. WorkflowAgentOutcome,
  21. WorkflowMeta,
  22. WorkflowPhase,
  23. WorkflowResult,
  24. WorkflowResultInfo,
  25. WorkflowRun,
  26. WorkflowRunInfo,
  27. WorkflowStartRequest,
  28. WorkflowStopReason,
  29. } from './types.ts'
  30. declare module 'cordis' {
  31. interface Context {
  32. workflows: WorkflowService
  33. }
  34. interface Events {
  35. /**
  36. * A workflow run started — the script's meta block validated, the body
  37. * about to execute. Paired with {@link Events['workflow/end']}.
  38. * @param info - the run's identity snapshot (id + meta).
  39. * @mode emit
  40. */
  41. 'workflow/start'(info: WorkflowRunInfo): void
  42. /**
  43. * The script entered a phase (a `phase(title)` call) — progress grouping
  44. * for observers; no execution semantics.
  45. * @param info - the run's identity snapshot.
  46. * @param title - the phase title, verbatim.
  47. * @mode emit
  48. */
  49. 'workflow/phase'(info: WorkflowRunInfo, title: string): void
  50. /**
  51. * The script emitted a narration line (a `log(message)` call).
  52. * @param info - the run's identity snapshot.
  53. * @param message - the logged message, verbatim.
  54. * @mode emit
  55. */
  56. 'workflow/log'(info: WorkflowRunInfo, message: string): void
  57. /**
  58. * One `agent()` call established a published child run. Paired with
  59. * {@link Events['workflow/agent-end']} by `agent.seq`. A call that never
  60. * receives a published run from the provider emits neither
  61. * event in this pair.
  62. * @param info - the run's identity snapshot.
  63. * @param agent - the call's sequence number, label, phase, and child id.
  64. * @mode emit
  65. */
  66. 'workflow/agent-start'(info: WorkflowRunInfo, agent: WorkflowAgentInfo): void
  67. /**
  68. * One `agent()` call settled (clean result, child failure, or run
  69. * cancellation). Paired with {@link Events['workflow/agent-start']} by
  70. * `agent.seq`, exactly once per started call on every stop path — on an
  71. * engine termination path (a worker killed past its grace) the end is
  72. * engine-synthesized with outcome `'cancelled'`.
  73. * @param info - the run's identity snapshot.
  74. * @param agent - the call identity plus its outcome.
  75. * @mode emit
  76. */
  77. 'workflow/agent-end'(info: WorkflowRunInfo, agent: WorkflowAgentEndInfo): void
  78. /**
  79. * A workflow run settled (any stop reason). Fired when
  80. * {@link WorkflowRun.result} resolves. Paired with
  81. * {@link Events['workflow/start']}.
  82. * @param info - the run's identity snapshot.
  83. * @param result - the outcome data (stop reason, error, agent count) —
  84. * deliberately WITHOUT the result value (see {@link WorkflowResultInfo}).
  85. * @mode emit
  86. */
  87. 'workflow/end'(info: WorkflowRunInfo, result: WorkflowResultInfo): void
  88. }
  89. }
  90. /** The full set of `workflow/*` event names {@link WorkflowService.emitWorkflowEvent} dispatches. */
  91. export type WorkflowEventName =
  92. | 'workflow/start'
  93. | 'workflow/phase'
  94. | 'workflow/log'
  95. | 'workflow/agent-start'
  96. | 'workflow/agent-end'
  97. | 'workflow/end'
  98. /**
  99. * Machine-routable fatal workflow failures: parse/meta/argument/schema errors,
  100. * resource caps, subagent infrastructure failures, unserializable boundary
  101. * values, and cancellation. An ordinary child failure resolves its item to
  102. * `null` and is not one of these fatal codes.
  103. */
  104. export type WorkflowErrorCode =
  105. | 'SCRIPT_PARSE'
  106. | 'META_INVALID'
  107. | 'INVALID_ARGUMENT'
  108. | 'UNSUPPORTED_OPTION'
  109. | 'UNSUPPORTED_SCHEMA'
  110. | 'AGENT_CAP'
  111. | 'ITEM_CAP'
  112. | 'AGENT_START'
  113. | 'AGENT_RESULT'
  114. | 'RESULT_UNSERIALIZABLE'
  115. | 'CANCELLED'
  116. /**
  117. * Typed error for workflow-seam failures. Extends {@link HarnessError}, so the
  118. * `code` is machine-routable taxonomy. `fatal` drives the combinator
  119. * discipline: `parallel()`/`pipeline()` re-throw a fatal error (a typo'd
  120. * option or a tripped cap must kill the script loudly), and reserve the
  121. * per-item `null` for child-run failures and ordinary in-stage script errors.
  122. * Every {@link WorkflowErrorCode} is fatal; the flag exists so the
  123. * distinction is explicit at every catch site rather than implied.
  124. */
  125. export class WorkflowError extends HarnessError {
  126. /** Whether combinators must propagate this error instead of nulling the item. */
  127. readonly fatal: boolean
  128. constructor(message: string, code: WorkflowErrorCode, options?: ErrorOptions & { fatal?: boolean }) {
  129. super(message, code, options)
  130. this.name = 'WorkflowError'
  131. this.fatal = options?.fatal ?? true
  132. }
  133. }
  134. /**
  135. * Whether combinators must re-throw `error` instead of mapping the item to `null`.
  136. * @param error - any thrown value; fatality is host `instanceof` (unforgeable from a script realm).
  137. * @returns true iff `error` is a {@link WorkflowError} whose `fatal` flag is set.
  138. */
  139. export function isFatalWorkflowError(error: unknown): boolean {
  140. return error instanceof WorkflowError && error.fatal
  141. }
  142. /**
  143. * Workflow Service Definition contract. Invalid requests throw before publication; a live
  144. * run is holder-owned, its result never rejects, cancellation and disposal are
  145. * bounded, and disposal waits for child cleanup within that bound. Lifecycle
  146. * listener failures are contained, and `workflow/end` fires exactly once as the
  147. * result settles.
  148. */
  149. export abstract class WorkflowService extends Service {
  150. constructor(ctx: Context) {
  151. super(ctx, 'workflows')
  152. }
  153. /**
  154. * Parse and execute a workflow script.
  155. * @param request - the script, its `args`, the parent agent, and an
  156. * optional cancel signal.
  157. * @returns the live run; its `result` resolves when the script settles.
  158. */
  159. abstract start(request: WorkflowStartRequest): WorkflowRun
  160. /**
  161. * Emit a lifecycle event while containing and logging each listener failure.
  162. * @param name - the `workflow/*` event to dispatch.
  163. * @param args - the event's payload, matching its declared signature.
  164. */
  165. protected emitWorkflowEvent(name: WorkflowEventName, ...args: unknown[]): void {
  166. for (const callback of this.ctx.events.dispatch('emit', [name, ...args])) {
  167. try {
  168. const returned: unknown = (callback as (...payload: unknown[]) => unknown)(...args)
  169. void Promise.resolve(returned).catch((error: unknown) => {
  170. this.ctx.logger.warn(`workflow: ${name} listener rejected: ${renderListenerError(error)}`)
  171. })
  172. } catch (error: unknown) {
  173. this.ctx.logger.warn(`workflow: ${name} listener threw: ${renderListenerError(error)}`)
  174. }
  175. }
  176. }
  177. }
  178. /**
  179. * Render any thrown value without violating listener containment.
  180. * @param error - any thrown value.
  181. * @returns `String(error)`, or a fixed label when even coercion throws.
  182. */
  183. function renderListenerError(error: unknown): string {
  184. try {
  185. return String(error)
  186. } catch {
  187. // String coercion itself may throw.
  188. return '[unrenderable thrown value]'
  189. }
  190. }
  191. export default WorkflowService