import { assertValidAgentNode, assertValidActionNode, assertValidCheckpointNode, assertValidComputeNode, assertValidNotifyNode, assertValidShellActionNode, assertValidWorkflowDefinitionShape, } from "./schema.js"; import type { AgentNodeDefinition, AssistantAgentNodeDefinition, AssistantMessageOutput, SubmittedAgentNodeDefinition, ActionNodeDefinition, CheckpointNodeDefinition, ComputeNodeDefinition, FunctionActionNodeDefinition, NotifyNodeDefinition, ShellActionNodeDefinition, WorkflowDefinition, WorkflowExitMap, WorkflowIncludeDefinition, WorkflowIncludedResult, WorkflowIncludeMap, WorkflowInputOf, WorkflowNodeContext, WorkflowNodeDefinition, WorkflowTypedEdge, WorkflowValueParser, } from "./types.js"; const WORKFLOW_DEFINITION_BRAND = Symbol.for("pi-workflows.definition"); type WorkflowDefinitionInput< TInput, TNodes extends Record, TIncludes extends WorkflowIncludeMap, TExits extends WorkflowExitMap, > = Omit< WorkflowDefinition, "input" | "nodes" | "includes" | "exits" | "edges" > & { input?: WorkflowValueParser; nodes: TNodes; includes?: TIncludes; exits?: TExits; edges: WorkflowTypedEdge[]; }; export function defineWorkflow< TInput = unknown, const TNodes extends Record = Record< string, WorkflowNodeDefinition >, const TIncludes extends WorkflowIncludeMap = Record, const TExits extends WorkflowExitMap = Record, >( definition: WorkflowDefinitionInput, ): WorkflowDefinition & { nodes: TNodes; includes?: TIncludes; exits?: TExits; } { assertValidWorkflowDefinitionShape(definition as WorkflowDefinition); const typed = definition as WorkflowDefinition & { nodes: TNodes; includes?: TIncludes; exits?: TExits; }; if (isWorkflowDefinition(typed)) { return typed; } Object.defineProperty(typed, WORKFLOW_DEFINITION_BRAND, { value: true, enumerable: false, configurable: false, writable: false, }); return typed; } export function includeWorkflow>( workflow: TWorkflow, options?: { input?: ( context: WorkflowNodeContext, ) => Promise> | WorkflowInputOf; }, ): WorkflowIncludeDefinition; export function includeWorkflow>( definition: WorkflowIncludeDefinition, ): WorkflowIncludeDefinition; export function includeWorkflow>( workflowOrDefinition: TWorkflow | WorkflowIncludeDefinition, options: { input?: ( context: WorkflowNodeContext, ) => Promise> | WorkflowInputOf; } = {}, ): WorkflowIncludeDefinition { const definition = isWorkflowDefinition(workflowOrDefinition) ? { workflow: workflowOrDefinition, ...options } : workflowOrDefinition; if (typeof definition.workflow !== "string" && !isWorkflowDefinition(definition.workflow)) { throw new Error("Included workflow must be a defined workflow or a workflow reference"); } if (definition.input !== undefined && typeof definition.input !== "function") { throw new Error("Included workflow input must be a function"); } if (definition.contract !== undefined && !isWorkflowDefinition(definition.contract)) { throw new Error("Included workflow contract must be defined with defineWorkflow"); } return definition; } export function includedResult>( workflow: TWorkflow, value: unknown, ): WorkflowIncludedResult { if (value === null || typeof value !== "object" || Array.isArray(value)) { throw new Error(`Included ${workflow.name} result must be an object`); } const exit = (value as { exit?: unknown }).exit; if (typeof exit !== "string" || !Object.hasOwn(workflow.exits ?? {}, exit)) { throw new Error(`Included ${workflow.name} result has unknown exit ${JSON.stringify(exit)}`); } if (!Object.hasOwn(value, "output")) { throw new Error(`Included ${workflow.name} result requires output`); } return value as WorkflowIncludedResult; } /** Preserve exact workflow types while checking duplicate registry names. */ export function defineWorkflowRegistry< const TRegistry extends Record>, >(registry: TRegistry): Readonly { const names = new Set(); for (const [key, workflow] of Object.entries(registry)) { if (!isWorkflowDefinition(workflow)) { throw new Error(`Workflow registry entry ${key} is not defined with defineWorkflow`); } if (names.has(workflow.name)) { throw new Error(`Workflow registry has duplicate workflow name: ${workflow.name}`); } names.add(workflow.name); } return Object.freeze({ ...registry }); } export function isWorkflowDefinition(value: unknown): value is WorkflowDefinition { return ( value != null && typeof value === "object" && (value as Record)[WORKFLOW_DEFINITION_BRAND] === true ); } export function assistantMessage(options: { maxChars?: number } = {}): AssistantMessageOutput { if (options === null || typeof options !== "object" || Array.isArray(options)) { throw new Error("Invalid workflow definition: assistantMessage options must be an object"); } const unknown = Object.keys(options).filter((key) => key !== "maxChars"); if (unknown.length > 0) { throw new Error( `Invalid workflow definition: assistantMessage has unknown option ${JSON.stringify(unknown[0])}`, ); } if ( options.maxChars !== undefined && (!Number.isInteger(options.maxChars) || options.maxChars <= 0) ) { throw new Error( "Invalid workflow definition: assistantMessage maxChars must be a positive integer", ); } return Object.freeze({ kind: "assistant-message" as const, ...(options.maxChars !== undefined ? { maxChars: options.maxChars } : {}), }); } export function agent( definition: Omit, ): SubmittedAgentNodeDefinition; export function agent( definition: Omit, ): AssistantAgentNodeDefinition; export function agent(definition: Omit): AgentNodeDefinition { const node = { nodeType: "agent" as const, ...definition, } as AgentNodeDefinition; assertValidAgentNode(node); return node; } export function compute( definition: Omit, ): ComputeNodeDefinition { const node: ComputeNodeDefinition = { nodeType: "compute", ...definition, }; assertValidComputeNode(node); return node; } export function notify(definition: Omit): NotifyNodeDefinition { const node: NotifyNodeDefinition = { nodeType: "notify", ...definition, }; assertValidNotifyNode(node); return node; } export function action( definition: Omit, ): FunctionActionNodeDefinition; export function action( definition: Omit, ): ShellActionNodeDefinition; export function action( definition: | Omit | Omit, ): ActionNodeDefinition { const node: ActionNodeDefinition = { nodeType: "action", ...definition, }; assertValidActionNode(node); return node; } export function shell( definition: Omit, ): ShellActionNodeDefinition { const node: ShellActionNodeDefinition = { nodeType: "action", ...definition, }; assertValidShellActionNode(node); return node; } export function checkpoint( definition: Omit = {}, ): CheckpointNodeDefinition { const node: CheckpointNodeDefinition = { nodeType: "checkpoint", ...definition, }; assertValidCheckpointNode(node); return node; }