import { apiRequest, apiRequestRaw, getAppBaseUrl } from './common.ts' import type { Variable } from './canvas.ts' import { getWorkflowModuleRule } from './workflow-module-rules.ts' export type WorkflowVersion = 'v3' export type WorkflowActionId = | 'agent' | 'canvas' | 'code-execution' | 'condition' | 'document-cropper' | 'document-splitter' | 'document-templater' | 'llm-completion' | 'map' | 'stop' | 'template' export const WORKFLOW_ACTION_IDS_V3_UI: WorkflowActionId[] = [ 'agent', 'canvas', 'code-execution', 'condition', 'document-cropper', 'document-splitter', 'document-templater', 'llm-completion', 'map', 'stop', 'template', ] export const WORKFLOW_ACTION_IDS_SHARED: WorkflowActionId[] = [...WORKFLOW_ACTION_IDS_V3_UI] export type WorkflowPortDirection = 'in' | 'out' export interface WorkflowPort { id: string dir: string type: string } export interface WorkflowNode { id: string name: string actionId: string input?: Record | null ports?: WorkflowPort[] [key: string]: unknown } export interface WorkflowEdgeEndpoint { id: string portId: string | null } export interface WorkflowEdge { id: string source: WorkflowEdgeEndpoint target: WorkflowEdgeEndpoint } export interface WorkflowGraph { nodes: WorkflowNode[] edges: WorkflowEdge[] } export interface WorkflowGraphValidationIssue { code: string message: string nodeId?: string edgeId?: string path?: string } export interface WorkflowGraphValidationResult { valid: boolean errors: WorkflowGraphValidationIssue[] warnings: WorkflowGraphValidationIssue[] } export interface ValidateWorkflowGraphOptions { version?: WorkflowVersion /** * When true, enforce the node catalog used by each UI route. */ strictUiCompatibility?: boolean /** * Deprecated no-op retained for older callers. Workflows always use v3 semantics. */ strictReferenceCompatibility?: boolean /** * When true, enforce graph shapes that are authorable through the Tela UI. */ strictAuthoring?: boolean } export interface WorkflowPromptVersion { id: string promptId: string title: string content: string | null markdownContent: string | null configuration: Record | null variables: Variable[] promoted: boolean draft: boolean isWorkflow: boolean | null workflowSpec: string | null graph: string | null createdAt: string updatedAt: string [key: string]: unknown } export interface CreateWorkflowVersionPayload { promptId: string title: string workflowSpec: Record | string graph: WorkflowGraph | string variables?: Variable[] configuration?: Record | null content?: string | null markdownContent?: string | null promoted?: boolean draft?: boolean version?: WorkflowVersion strictUiCompatibility?: boolean strictReferenceCompatibility?: boolean strictAuthoring?: boolean } export interface UpdateWorkflowVersionPayload { title?: string workflowSpec?: Record | string graph?: WorkflowGraph | string variables?: Variable[] configuration?: Record | null content?: string | null markdownContent?: string | null promoted?: boolean draft?: boolean version?: WorkflowVersion strictUiCompatibility?: boolean strictReferenceCompatibility?: boolean strictAuthoring?: boolean } export interface UpdateWorkflowVersionPreservingStatePayload extends UpdateWorkflowVersionPayload { promptId?: string } export type WorkflowEnvironment = 'test-case' | 'craft-preview' | 'api' | 'workstation' export interface WorkflowDefinition { inputs?: Record outputs?: Record steps?: Record | unknown[] settings?: { parallelExecution?: boolean [key: string]: unknown } [key: string]: unknown } export interface RunWorkflowPayload { id?: string definition: WorkflowDefinition workflowSpec?: Record inputs: Record promptVersionId: string promptId?: string graph?: WorkflowGraph | null environment?: WorkflowEnvironment startStepId?: string | null endStepId?: string | null cacheStrategy?: { enabled: boolean ttl?: number } completionId?: string customTags?: string[] allowedCredentials?: string[] version?: WorkflowVersion strictUiCompatibility?: boolean strictReferenceCompatibility?: boolean strictAuthoring?: boolean } export type WorkflowRunStatus = | 'pending' | 'running' | 'completed' | 'failed' | 'requires-action' | 'cancelled' | 'timeout' | 'waiting-for-review' export interface WorkflowRun { id: string promptId: string promptVersionId: string status: WorkflowRunStatus input?: Record output?: unknown graph?: WorkflowGraph steps?: Array> startedAt?: string | null completedAt?: string | null createdAt?: string updatedAt?: string [key: string]: unknown } export interface WaitForWorkflowRunOptions { timeoutMs?: number pollIntervalMs?: number terminalStatuses?: WorkflowRunStatus[] } export interface ListWorkflowRunsOptions { promptId: string promptVersionId?: string limit?: number status?: 'pending' | 'completed' | 'failed' | 'cancelled' offset?: number orderBy?: 'createdAt' | 'updatedAt' | 'completedAt' orderDirection?: 'asc' | 'desc' } export interface WorkflowVariableDefinition { id?: string name: string type: 'file' | 'text' required: boolean description?: string processingOptions?: { allowMultimodal: boolean } } export interface WorkflowAvailableVariable { id?: string key?: string label?: string type?: string source?: string order?: number schema?: unknown value?: unknown [key: string]: unknown } export interface GetWorkflowVariablesPayload { graph?: WorkflowGraph | string promptVersionId?: string input?: WorkflowVariableDefinition[] version?: WorkflowVersion strictUiCompatibility?: boolean strictReferenceCompatibility?: boolean strictAuthoring?: boolean } export interface WorkflowVariablesResponse { variables: WorkflowAvailableVariable[] stepId?: string } export interface GetWorkflowAttributesOptions { environment: WorkflowEnvironment promptVersionId: string } const DEFAULT_TERMINAL_STATUSES: WorkflowRunStatus[] = ['completed', 'failed', 'cancelled', 'timeout'] const CONTROL_FLOW_ACTIONS = new Set(['condition', 'map', 'stop']) function normalizeVersion(_version?: WorkflowVersion): WorkflowVersion { return 'v3' } function getAllowedActions(strictUiCompatibility: boolean): Set { if (!strictUiCompatibility) return new Set(WORKFLOW_ACTION_IDS_SHARED) return new Set(WORKFLOW_ACTION_IDS_V3_UI) } function parseGraphInput(graph: WorkflowGraph | string): WorkflowGraph { if (typeof graph === 'string') return JSON.parse(graph) as WorkflowGraph return graph } function deriveVariablesFromSpec(spec: Record | string): Variable[] { const parsed = typeof spec === 'string' ? JSON.parse(spec) : spec const inputs = parsed?.inputs as Record> | undefined if (!inputs) return [] return Object.values(inputs).map((input) => { const variable: Variable = { id: input.id as string | undefined, name: input.name as string, type: input.type as 'file' | 'text', required: (input.required as boolean) ?? true, } if (input.description) variable.description = input.description as string const processingOptions = input.processingOptions as { allowMultimodal?: boolean } | undefined if (processingOptions) variable.processingOptions = { allowMultimodal: processingOptions.allowMultimodal ?? false } return variable }) } type WorkflowValidationContext = Pick< ValidateWorkflowGraphOptions, 'version' | 'strictUiCompatibility' | 'strictReferenceCompatibility' | 'strictAuthoring' > function resolveGraphWithValidation( graph: WorkflowGraph | string, options: WorkflowValidationContext = {}, ): WorkflowGraph { const graphObject = parseGraphInput(graph) assertWorkflowGraphValid(graphObject, { version: normalizeVersion(options.version), strictUiCompatibility: options.strictUiCompatibility, strictReferenceCompatibility: options.strictReferenceCompatibility, strictAuthoring: options.strictAuthoring, }) return graphObject } function resolveOptionalGraphWithValidation( graph: WorkflowGraph | string | undefined, options: WorkflowValidationContext = {}, ): WorkflowGraph | undefined { if (graph === undefined) return undefined return resolveGraphWithValidation(graph, options) } function stringifyJsonInput(value: Record | string): string { if (typeof value === 'string') return value return JSON.stringify(value) } function hasPort(node: WorkflowNode, direction: WorkflowPortDirection, type: string): boolean { return Boolean(node.ports?.some(port => port.dir === direction && port.type === type)) } function getPort(node: WorkflowNode, portId: string | null): WorkflowPort | undefined { if (!portId) return undefined return node.ports?.find(port => port.id === portId) } function isEdgeEndpoint(value: unknown): value is WorkflowEdgeEndpoint { return Boolean( value && typeof value === 'object' && typeof (value as Partial).id === 'string' && (value as Partial).id!.length > 0, ) } function isWorkflowEdgeObject(value: unknown): value is WorkflowEdge { return Boolean(value && typeof value === 'object' && !Array.isArray(value)) } function isRecord(value: unknown): value is Record { return Boolean(value && typeof value === 'object' && !Array.isArray(value)) } function edgeKey(endpoint: WorkflowEdgeEndpoint): string { return `${endpoint.id}:${endpoint.portId ?? ''}` } function findCycle(nodes: WorkflowNode[], edges: WorkflowEdge[]): string[] | null { const outgoing = new Map() for (const node of nodes) outgoing.set(node.id, []) for (const edge of edges) outgoing.get(edge.source.id)?.push(edge.target.id) const visiting = new Set() const visited = new Set() const stack: string[] = [] function visit(id: string): string[] | null { if (visiting.has(id)) { const start = stack.indexOf(id) return stack.slice(start).concat(id) } if (visited.has(id)) return null visiting.add(id) stack.push(id) for (const child of outgoing.get(id) ?? []) { const cycle = visit(child) if (cycle) return cycle } stack.pop() visiting.delete(id) visited.add(id) return null } for (const node of nodes) { const cycle = visit(node.id) if (cycle) return cycle } return null } function collectReachable(rootIds: string[], edges: WorkflowEdge[]): Set { const reachable = new Set() const queue = [...rootIds] while (queue.length > 0) { const id = queue.shift()! if (reachable.has(id)) continue reachable.add(id) for (const edge of edges) { if (edge.source.id === id) queue.push(edge.target.id) } } return reachable } interface ConditionBranchScope { conditionNodeId: string branchPortId: string } function findConditionBranchScope( nodeId: string, incoming: Map, nodeMap: Map, visited = new Set(), ): ConditionBranchScope | null { if (visited.has(nodeId)) return null visited.add(nodeId) for (const edge of incoming.get(nodeId) ?? []) { const sourceNode = nodeMap.get(edge.source.id) const sourcePort = sourceNode ? getPort(sourceNode, edge.source.portId) : undefined if (sourceNode && sourcePort?.type === 'branch') { return { conditionNodeId: sourceNode.id, branchPortId: sourcePort.id, } } const scope = findConditionBranchScope(edge.source.id, incoming, nodeMap, visited) if (scope) return scope } return null } function getConditionScopedIncoming( edge: WorkflowEdge, incoming: Map, nodeMap: Map, ): ConditionBranchScope | null { const sourceNode = nodeMap.get(edge.source.id) const sourcePort = sourceNode ? getPort(sourceNode, edge.source.portId) : undefined if (sourceNode && sourcePort?.type === 'branch') { return { conditionNodeId: sourceNode.id, branchPortId: sourcePort.id, } } return findConditionBranchScope(edge.source.id, incoming, nodeMap) } function isConditionScopedIncoming( edge: WorkflowEdge, incoming: Map, nodeMap: Map, ): boolean { return getConditionScopedIncoming(edge, incoming, nodeMap) !== null } function validateAuthoringShape( nodes: WorkflowNode[], edges: WorkflowEdge[], nodeMap: Map, ): WorkflowGraphValidationIssue[] { const errors: WorkflowGraphValidationIssue[] = [] const validEdges = edges.filter(edge => isEdgeEndpoint(edge.source) && isEdgeEndpoint(edge.target)) const incoming = new Map() const outgoingBySourcePort = new Map() for (const node of nodes) incoming.set(node.id, []) for (const edge of validEdges) { incoming.get(edge.target.id)?.push(edge) const key = edgeKey(edge.source) const outgoing = outgoingBySourcePort.get(key) ?? [] outgoing.push(edge) outgoingBySourcePort.set(key, outgoing) } const roots = nodes.filter(node => (incoming.get(node.id) ?? []).length === 0) if (nodes.length > 0 && roots.length !== 1) { errors.push({ code: 'invalid_root_count', message: `UI-authored workflows must have exactly one main root. Found ${roots.length}.`, }) } if (roots.some(node => node.actionId === 'stop')) { const stopRoot = roots.find(node => node.actionId === 'stop')! errors.push({ code: 'stop_cannot_be_root', message: `Stop node '${stopRoot.id}' cannot be the main root in authoring mode.`, nodeId: stopRoot.id, }) } const cycle = findCycle(nodes, validEdges) if (cycle) { errors.push({ code: 'cycle_detected', message: `Workflow graph contains a cycle: ${cycle.join(' -> ')}.`, }) } if (roots.length === 1) { const reachable = collectReachable([roots[0]!.id], validEdges) for (const node of nodes) { if (!reachable.has(node.id)) { errors.push({ code: 'unreachable_node', message: `Node '${node.id}' is not reachable from the main root.`, nodeId: node.id, }) } } } for (const [sourceKey, sourceEdges] of outgoingBySourcePort.entries()) { const [nodeId, portId] = sourceKey.split(':') const node = nodeMap.get(nodeId!) const port = node?.ports?.find(candidate => candidate.id === portId) if (!node || !port) continue if (port.type === 'default' && sourceEdges.length > 1) { errors.push({ code: 'default_fan_out', message: `Node '${node.id}' default output has ${sourceEdges.length} outgoing edges; UI-authored default flow supports at most one.`, nodeId: node.id, }) } } for (const node of nodes) { if (node.actionId !== 'condition') continue const branchPorts = (node.ports ?? []).filter(port => port.dir === 'out' && port.type === 'branch') for (const port of branchPorts) { const count = outgoingBySourcePort.get(`${node.id}:${port.id}`)?.length ?? 0 if (count !== 1) { errors.push({ code: 'condition_branch_edge_count', message: `Condition node '${node.id}' branch port '${port.id}' must have exactly one outgoing edge in authoring mode. Found ${count}.`, nodeId: node.id, }) } } } for (const node of nodes) { const incomingEdges = incoming.get(node.id) ?? [] if (incomingEdges.length <= 1) continue const incomingScopes = incomingEdges.map(edge => getConditionScopedIncoming(edge, incoming, nodeMap), ) const allConditionScoped = incomingScopes.every(scope => scope !== null) if (!allConditionScoped) { errors.push({ code: 'arbitrary_fan_in', message: `Node '${node.id}' has ${incomingEdges.length} incoming edges that are not all condition-scoped.`, nodeId: node.id, }) } else { const conditionIds = new Set( incomingScopes.map(scope => scope.conditionNodeId), ) if (conditionIds.size > 1) { errors.push({ code: 'condition_connect_scope_mismatch', message: `Node '${node.id}' has connection fan-in from multiple condition scopes; UI-authored connect targets must merge branches from one condition.`, nodeId: node.id, }) } } } for (const node of nodes) { if (node.actionId !== 'stop') continue const incomingEdges = incoming.get(node.id) ?? [] const fromConditionBranch = incomingEdges.some(edge => isConditionScopedIncoming(edge, incoming, nodeMap), ) if (!fromConditionBranch) { errors.push({ code: 'stop_not_at_condition_branch_end', message: `Stop node '${node.id}' must be placed at a condition branch end in authoring mode.`, nodeId: node.id, }) } } return errors } export function validateWorkflowGraph( graph: WorkflowGraph, options: ValidateWorkflowGraphOptions = {}, ): WorkflowGraphValidationResult { const errors: WorkflowGraphValidationIssue[] = [] const warnings: WorkflowGraphValidationIssue[] = [] const strictUiCompatibility = options.strictUiCompatibility ?? true if (!graph || typeof graph !== 'object') { return { valid: false, errors: [{ code: 'invalid_graph', message: 'Graph must be an object with nodes and edges.' }], warnings, } } if (!Array.isArray(graph.nodes)) { errors.push({ code: 'invalid_nodes', message: 'Graph.nodes must be an array.' }) } if (!Array.isArray(graph.edges)) { errors.push({ code: 'invalid_edges', message: 'Graph.edges must be an array.' }) } if (errors.length > 0) { return { valid: false, errors, warnings } } const nodes = graph.nodes const edges = graph.edges const allowedActions = getAllowedActions(strictUiCompatibility) const nodeIds = new Set() const edgeIds = new Set() const nodeMap = new Map() for (const node of nodes) { if (!node.id) { errors.push({ code: 'node_missing_id', message: 'Node is missing id.' }) continue } if (nodeIds.has(node.id)) { errors.push({ code: 'duplicate_node_id', message: `Duplicate node id '${node.id}'.`, nodeId: node.id, }) } nodeIds.add(node.id) nodeMap.set(node.id, node) if (!node.actionId) { errors.push({ code: 'node_missing_action', message: `Node '${node.id}' is missing actionId.`, nodeId: node.id, }) continue } if (!allowedActions.has(node.actionId)) { errors.push({ code: 'unsupported_action_for_version', message: `Node '${node.id}' uses action '${node.actionId}' not allowed for workflow v3.`, nodeId: node.id, }) } if (!Array.isArray(node.ports) || node.ports.length === 0) { errors.push({ code: 'node_missing_ports', message: `Node '${node.id}' must define ports.`, nodeId: node.id, }) continue } const portIds = new Set() for (const port of node.ports) { if (!port.id) { errors.push({ code: 'port_missing_id', message: `Node '${node.id}' has a port without id.`, nodeId: node.id, }) continue } if (portIds.has(port.id)) { errors.push({ code: 'duplicate_port_id', message: `Node '${node.id}' has duplicate port id '${port.id}'.`, nodeId: node.id, }) } portIds.add(port.id) if (port.dir !== 'in' && port.dir !== 'out') { errors.push({ code: 'invalid_port_direction', message: `Node '${node.id}' port '${port.id}' has invalid direction '${port.dir}'.`, nodeId: node.id, }) } } if (!CONTROL_FLOW_ACTIONS.has(node.actionId)) { if (!hasPort(node, 'in', 'default')) { errors.push({ code: 'missing_default_in', message: `Node '${node.id}' must have a default input port.`, nodeId: node.id, }) } if (!hasPort(node, 'out', 'default')) { errors.push({ code: 'missing_default_out', message: `Node '${node.id}' must have a default output port.`, nodeId: node.id, }) } } if (node.actionId === 'condition') { if (!hasPort(node, 'in', 'default')) { errors.push({ code: 'condition_missing_default_in', message: `Condition node '${node.id}' must have a default input port.`, nodeId: node.id, }) } const branchCount = node.ports.filter(port => port.dir === 'out' && port.type === 'branch').length if (branchCount === 0) { errors.push({ code: 'condition_missing_branch_ports', message: `Condition node '${node.id}' must have at least one branch output port.`, nodeId: node.id, }) } const input = isRecord(node.input) ? node.input : undefined const expectedBranchPortIds = new Set() const seenInputBranchPortIds = new Set() const cases = Array.isArray(input?.cases) ? input.cases : [] for (const [index, workflowCase] of cases.entries()) { if (!isRecord(workflowCase)) { continue } const portId = workflowCase.portId if (typeof portId === 'string' && portId.length > 0) { if (seenInputBranchPortIds.has(portId)) { errors.push({ code: 'condition_branch_port_duplicate', message: `Condition node '${node.id}' has duplicate input branch port id '${portId}'.`, nodeId: node.id, path: `input.cases[${index}].portId`, }) } seenInputBranchPortIds.add(portId) expectedBranchPortIds.add(portId) } } const defaultBranch = isRecord(input?.default) ? input.default : undefined const defaultPortId = defaultBranch?.portId if (typeof defaultPortId === 'string' && defaultPortId.length > 0) { if (seenInputBranchPortIds.has(defaultPortId)) { errors.push({ code: 'condition_branch_port_duplicate', message: `Condition node '${node.id}' has duplicate default branch port id '${defaultPortId}'.`, nodeId: node.id, path: 'input.default.portId', }) } seenInputBranchPortIds.add(defaultPortId) expectedBranchPortIds.add(defaultPortId) } const actualBranchPortIds = new Set( node.ports .filter(port => port.dir === 'out' && port.type === 'branch') .map(port => port.id), ) for (const portId of expectedBranchPortIds) { if (!actualBranchPortIds.has(portId)) { errors.push({ code: 'condition_branch_port_missing', message: `Condition node '${node.id}' input references missing branch port '${portId}'.`, nodeId: node.id, }) } } for (const portId of actualBranchPortIds) { if (!expectedBranchPortIds.has(portId)) { errors.push({ code: 'condition_branch_port_orphaned', message: `Condition node '${node.id}' has branch port '${portId}' not present in input cases/default.`, nodeId: node.id, }) } } } if (node.actionId === 'map') { if (!hasPort(node, 'in', 'default')) { errors.push({ code: 'map_missing_default_in', message: `Map node '${node.id}' must have a default input port.`, nodeId: node.id, }) } if (!hasPort(node, 'out', 'subgraph')) { errors.push({ code: 'map_missing_subgraph_out', message: `Map node '${node.id}' must have a subgraph output port.`, nodeId: node.id, }) } if (!hasPort(node, 'out', 'default')) { errors.push({ code: 'map_missing_default_out', message: `Map node '${node.id}' must have a default continuation output port.`, nodeId: node.id, }) } if (node.input && typeof node.input === 'object' && 'steps' in node.input) { warnings.push({ code: 'map_steps_legacy_warning', message: `Map node '${node.id}' contains input.steps, which is legacy workflow state and ignored in v3 graph-driven execution.`, nodeId: node.id, }) } } if (node.actionId === 'stop') { if (!hasPort(node, 'in', 'default')) { errors.push({ code: 'stop_missing_default_in', message: `Stop node '${node.id}' must have a default input port.`, nodeId: node.id, }) } } const moduleRule = getWorkflowModuleRule(node.actionId) if (moduleRule) { errors.push(...moduleRule.validate(node)) } } const edgeObjects: WorkflowEdge[] = [] for (const [index, edge] of edges.entries()) { if (!isWorkflowEdgeObject(edge)) { errors.push({ code: 'invalid_edge', message: `Graph.edges[${index}] must be an edge object.`, path: `edges[${index}]`, }) continue } edgeObjects.push(edge) if (!edge.id) { errors.push({ code: 'edge_missing_id', message: 'Edge is missing id.' }) continue } if (edgeIds.has(edge.id)) { errors.push({ code: 'duplicate_edge_id', message: `Duplicate edge id '${edge.id}'.`, edgeId: edge.id, }) } edgeIds.add(edge.id) const sourceEndpointValid = isEdgeEndpoint(edge.source) const targetEndpointValid = isEdgeEndpoint(edge.target) if (!sourceEndpointValid) { errors.push({ code: 'edge_source_missing', message: `Edge '${edge.id}' is missing a source endpoint.`, edgeId: edge.id, }) } if (!targetEndpointValid) { errors.push({ code: 'edge_target_missing', message: `Edge '${edge.id}' is missing a target endpoint.`, edgeId: edge.id, }) } if (!sourceEndpointValid || !targetEndpointValid) continue const sourceNode = nodeMap.get(edge.source.id) const targetNode = nodeMap.get(edge.target.id) if (!sourceNode) { errors.push({ code: 'edge_source_node_missing', message: `Edge '${edge.id}' references missing source node '${edge.source.id}'.`, edgeId: edge.id, }) continue } if (!targetNode) { errors.push({ code: 'edge_target_node_missing', message: `Edge '${edge.id}' references missing target node '${edge.target.id}'.`, edgeId: edge.id, }) continue } if (!edge.source?.portId) { errors.push({ code: 'edge_source_port_missing', message: `Edge '${edge.id}' is missing source portId.`, edgeId: edge.id, }) continue } if (!edge.target?.portId) { errors.push({ code: 'edge_target_port_missing', message: `Edge '${edge.id}' is missing target portId.`, edgeId: edge.id, }) continue } const sourcePort = getPort(sourceNode, edge.source.portId) const targetPort = getPort(targetNode, edge.target.portId) if (!sourcePort) { errors.push({ code: 'edge_source_port_not_found', message: `Edge '${edge.id}' source port '${edge.source.portId}' not found on node '${sourceNode.id}'.`, edgeId: edge.id, }) } else if (sourcePort.dir !== 'out') { errors.push({ code: 'edge_source_port_wrong_direction', message: `Edge '${edge.id}' source port '${sourcePort.id}' on node '${sourceNode.id}' must be 'out'.`, edgeId: edge.id, }) } if (!targetPort) { errors.push({ code: 'edge_target_port_not_found', message: `Edge '${edge.id}' target port '${edge.target.portId}' not found on node '${targetNode.id}'.`, edgeId: edge.id, }) } else if (targetPort.dir !== 'in') { errors.push({ code: 'edge_target_port_wrong_direction', message: `Edge '${edge.id}' target port '${targetPort.id}' on node '${targetNode.id}' must be 'in'.`, edgeId: edge.id, }) } } for (const node of nodes) { if (node.actionId === 'map') { const subgraphPortIds = new Set((node.ports ?? []) .filter(port => port.dir === 'out' && port.type === 'subgraph') .map(port => port.id)) const subgraphEdges = edgeObjects.filter(edge => isEdgeEndpoint(edge.source) && edge.source.id === node.id && Boolean(edge.source.portId) && subgraphPortIds.has(edge.source.portId as string), ) if (subgraphEdges.length === 0) { errors.push({ code: 'map_missing_subgraph_edge', message: `Map node '${node.id}' has no edge from its subgraph port.`, nodeId: node.id, }) } else if ((options.strictAuthoring ?? true) && subgraphEdges.length !== 1) { errors.push({ code: 'map_subgraph_root_count', message: `Map node '${node.id}' must have exactly one subgraph root edge in authoring mode. Found ${subgraphEdges.length}.`, nodeId: node.id, }) } } if (node.actionId === 'stop') { const hasOutgoingEdge = edgeObjects.some(edge => isEdgeEndpoint(edge.source) && edge.source.id === node.id) if (hasOutgoingEdge) { errors.push({ code: 'stop_has_outgoing_edges', message: `Stop node '${node.id}' should not have outgoing edges.`, nodeId: node.id, }) } } } if (options.strictAuthoring ?? true) { errors.push(...validateAuthoringShape(nodes, edgeObjects, nodeMap)) } return { valid: errors.length === 0, errors, warnings, } } export function assertWorkflowGraphValid( graph: WorkflowGraph, options: ValidateWorkflowGraphOptions = {}, ): void { const result = validateWorkflowGraph(graph, options) if (result.valid) return const details = result.errors .map(issue => `${issue.code}: ${issue.message}`) .join('\n') throw new Error(`Workflow graph validation failed:\n${details}`) } function assertWorkflowGraphValidLocallyWhenStrictAuthoring( graph: WorkflowGraph | undefined, options: ValidateWorkflowGraphOptions, ): void { if (!graph || options.strictAuthoring === false) return assertWorkflowGraphValid(graph, { version: options.version ?? 'v3', strictUiCompatibility: options.strictUiCompatibility, strictReferenceCompatibility: options.strictReferenceCompatibility, strictAuthoring: true, }) } export async function listPromptVersions(promptId: string): Promise { return await apiRequest(`/prompt-version?promptId=${promptId}`) } export async function getPromptVersion(versionId: string): Promise { return await apiRequest(`/prompt-version/${versionId}`) } export async function createWorkflowVersion(payload: CreateWorkflowVersionPayload): Promise { const graphObject = parseGraphInput(payload.graph) assertWorkflowGraphValidLocallyWhenStrictAuthoring(graphObject, { version: 'v3', strictUiCompatibility: payload.strictUiCompatibility, strictReferenceCompatibility: payload.strictReferenceCompatibility, strictAuthoring: payload.strictAuthoring, }) await assertWorkflowGraphValidServer(payload.promptId, graphObject) const variables = payload.variables ?? deriveVariablesFromSpec(payload.workflowSpec) return await apiRequest('/prompt-version', { method: 'POST', body: JSON.stringify({ promptId: payload.promptId, title: payload.title, workflowSpec: stringifyJsonInput(payload.workflowSpec), graph: JSON.stringify(graphObject), variables, configuration: payload.configuration !== undefined ? payload.configuration : { model: 'gpt-5', type: 'chat', temperature: 0 }, content: payload.content, markdownContent: payload.markdownContent, promoted: payload.promoted ?? false, draft: payload.draft ?? false, isWorkflow: true, }), }) } export async function createWorkflowVersionV3( payload: Omit, ): Promise { return await createWorkflowVersion(payload) } export async function updateWorkflowVersion( versionId: string, payload: UpdateWorkflowVersionPayload, ): Promise { const graphObject = payload.graph !== undefined ? parseGraphInput(payload.graph) : undefined if (graphObject !== undefined) { assertWorkflowGraphValidLocallyWhenStrictAuthoring(graphObject, { version: 'v3', strictUiCompatibility: payload.strictUiCompatibility, strictReferenceCompatibility: payload.strictReferenceCompatibility, strictAuthoring: payload.strictAuthoring, }) const existing = await getPromptVersion(versionId) await assertWorkflowGraphValidServer(existing.promptId, graphObject) } const graphString = graphObject ? JSON.stringify(graphObject) : undefined const variables = payload.variables ?? (payload.workflowSpec !== undefined ? deriveVariablesFromSpec(payload.workflowSpec) : undefined) return await apiRequest(`/prompt-version/${versionId}`, { method: 'PATCH', body: JSON.stringify({ title: payload.title, workflowSpec: payload.workflowSpec !== undefined ? stringifyJsonInput(payload.workflowSpec) : undefined, graph: graphString, variables, configuration: payload.configuration, content: payload.content, markdownContent: payload.markdownContent, promoted: payload.promoted, draft: payload.draft, isWorkflow: true, }), }) } export async function updateWorkflowVersionV3( versionId: string, payload: Omit, ): Promise { return await updateWorkflowVersion(versionId, payload) } export async function runWorkflow(payload: RunWorkflowPayload): Promise { if (payload.graph) { assertWorkflowGraphValidLocallyWhenStrictAuthoring(payload.graph, { version: 'v3', strictUiCompatibility: payload.strictUiCompatibility, strictReferenceCompatibility: payload.strictReferenceCompatibility, strictAuthoring: payload.strictAuthoring, }) const promptId = payload.promptId ?? (await getPromptVersion(payload.promptVersionId)).promptId await assertWorkflowGraphValidServer(promptId, payload.graph) } return await apiRequest('/workflow/test', { method: 'POST', body: JSON.stringify({ id: payload.id, definition: payload.definition, workflowSpec: payload.workflowSpec, inputs: payload.inputs, promptVersionId: payload.promptVersionId, promptId: payload.promptId, graph: payload.graph, environment: payload.environment, startStepId: payload.startStepId, endStepId: payload.endStepId, cacheStrategy: payload.cacheStrategy, completionId: payload.completionId, customTags: payload.customTags, allowedCredentials: payload.allowedCredentials, }), }) } export async function getWorkflowRun(workflowRunId: string): Promise { return await apiRequest(`/workflow/runs/${workflowRunId}`) } export async function listWorkflowRuns(options: ListWorkflowRunsOptions): Promise { const params = new URLSearchParams({ promptId: options.promptId }) if (options.promptVersionId) params.set('promptVersionId', options.promptVersionId) if (options.limit !== undefined) params.set('limit', String(options.limit)) if (options.status) params.set('status', options.status) if (options.offset !== undefined) params.set('offset', String(options.offset)) if (options.orderBy) params.set('orderBy', options.orderBy) if (options.orderDirection) params.set('orderDirection', options.orderDirection) return await apiRequest(`/workflow/runs?${params.toString()}`) } export async function cancelWorkflowRun(workflowRunId: string): Promise { return await apiRequest(`/workflow/runs/${workflowRunId}/cancel`, { method: 'POST', }) } export async function waitForWorkflowRun( workflowRunId: string, options: WaitForWorkflowRunOptions = {}, ): Promise { const timeoutMs = options.timeoutMs ?? 5 * 60 * 1000 const pollIntervalMs = options.pollIntervalMs ?? 1500 const terminalStatuses = options.terminalStatuses ?? DEFAULT_TERMINAL_STATUSES const startedAt = Date.now() while (Date.now() - startedAt < timeoutMs) { const run = await getWorkflowRun(workflowRunId) if (terminalStatuses.includes(run.status)) { return run } await new Promise(resolve => setTimeout(resolve, pollIntervalMs)) } throw new Error(`Workflow run ${workflowRunId} did not reach terminal status within ${timeoutMs}ms`) } async function requestWorkflowVariables( endpoint: string, payload: GetWorkflowVariablesPayload = {}, ): Promise { const graphObject = resolveOptionalGraphWithValidation(payload.graph, payload) return await apiRequest(endpoint, { method: 'POST', body: JSON.stringify({ graph: graphObject, promptVersionId: payload.promptVersionId, input: payload.input ?? [], }), }) } export async function getWorkflowVariables( promptId: string, payload: GetWorkflowVariablesPayload = {}, ): Promise { return await requestWorkflowVariables(`/workflow/${promptId}/variables`, payload) } export async function getWorkflowVariablesByStep( promptId: string, stepId: string, payload: GetWorkflowVariablesPayload = {}, ): Promise { return await requestWorkflowVariables(`/workflow/${promptId}/variables/${stepId}`, payload) } export async function getWorkflowVariablesFromVersion( promptId: string, versionId?: string, ): Promise { const query = versionId ? `?versionId=${versionId}` : '' return await apiRequest(`/workflow/${promptId}/variables${query}`) } export async function getWorkflowVariablesByStepFromVersion( promptId: string, stepId: string, versionId?: string, ): Promise { const query = versionId ? `?versionId=${versionId}` : '' return await apiRequest(`/workflow/${promptId}/variables/step/${stepId}${query}`) } export async function getWorkflowAttributes( promptId: string, options: GetWorkflowAttributesOptions, ): Promise> { const params = new URLSearchParams({ environment: options.environment, promptVersionId: options.promptVersionId, }) return await apiRequest>(`/workflow/${promptId}/attributes?${params.toString()}`) } export interface WorkflowActionCatalogEntry { actionId: string description: string version: number branching: boolean configSchema: Record outputSchema: Record references: Record | null publicMethods: Array<{ name: string, description: string }> } export interface WorkflowActionsCatalogResponse { actions: WorkflowActionCatalogEntry[] } export interface InferredVariablePathCompatibility { path: string formats: string[] defaultFormat: string } export type InferredVariableCompatibility = Record export interface InferredVariable { id: string name: string type: 'output' | 'input' scope: string order: number schema: Record | null metadata?: Record compatibility?: InferredVariableCompatibility [key: string]: unknown } export interface WorkflowDependenciesResponse { variables: InferredVariable[] } export interface WorkflowStepTypesResponse { dts: string } export type WorkflowServerValidationResult = | { valid: true } | { valid: false, errors: WorkflowGraphValidationIssue[] } export async function getWorkflowActions(): Promise { return await apiRequest('/workflow/actions') } export async function getWorkflowDependencies( promptId: string, versionId: string, ): Promise { const params = new URLSearchParams({ versionId }) return await apiRequest( `/workflow/${promptId}/dependencies?${params.toString()}`, ) } export async function getStepDependencies( promptId: string, versionId: string, stepId: string, ): Promise { const params = new URLSearchParams({ versionId }) return await apiRequest( `/workflow/${promptId}/dependencies/${stepId}?${params.toString()}`, ) } export async function getStepTypes( promptId: string, versionId: string, stepId: string, ): Promise { const params = new URLSearchParams({ versionId }) return await apiRequest( `/workflow/${promptId}/types/${stepId}?${params.toString()}`, ) } export async function validateWorkflowGraphServer( promptId: string, graph: WorkflowGraph, ): Promise { const response = await apiRequestRaw(`/workflow/${promptId}/validate`, { method: 'POST', body: { graph }, }) if (response.status === 200) { return { valid: true } } if (response.status === 400) { const body = await response.json() as { errors?: WorkflowGraphValidationIssue[] } return { valid: false, errors: body.errors ?? [] } } const error = await response.text() throw new Error(`validateWorkflowGraphServer failed: ${response.status} ${response.statusText} - ${error}`) } export async function assertWorkflowGraphValidServer( promptId: string, graph: WorkflowGraph, ): Promise { const result = await validateWorkflowGraphServer(promptId, graph) if (result.valid) return const details = result.errors.length > 0 ? result.errors.map(issue => `${issue.code}: ${issue.message}`).join('\n') : '(no error details returned by server)' throw new Error(`Workflow graph server validation failed:\n${details}`) } export function getWorkflowV3Url(promptId: string): string { return `${getAppBaseUrl()}/workflows/${promptId}` } export function getWorkflowUrl(promptId: string): string { return getWorkflowV3Url(promptId) } /** * Get the latest workflow version for a prompt (canvas) by ID. * Fetches all versions sorted by creation date and returns the most recent. */ export async function getLatestVersion(promptId: string): Promise { const versions = await listPromptVersions(promptId) if (versions.length === 0) { throw new Error(`No versions found for prompt '${promptId}'.`) } const sorted = versions.sort((a, b) => new Date(b.createdAt).getTime() - new Date(a.createdAt).getTime(), ) return sorted[0]! } /** * Fork an existing workflow version — returns a deep copy of the parsed graph and spec * ready for mutation, without modifying the original version. * * Usage: * const { graph, spec, configuration, version } = await forkVersion(versionId) * // mutate graph.nodes, spec, etc. * await createWorkflowVersion({ promptId, title: 'Patched workflow', workflowSpec: spec, graph, configuration }) */ export async function forkVersion(versionId: string): Promise<{ graph: WorkflowGraph spec: Record configuration: Record | null version: WorkflowPromptVersion }> { const version = await getPromptVersion(versionId) const graph: WorkflowGraph = version.graph ? JSON.parse(typeof version.graph === 'string' ? version.graph : JSON.stringify(version.graph)) : { nodes: [], edges: [] } const spec: Record = version.workflowSpec ? JSON.parse(typeof version.workflowSpec === 'string' ? version.workflowSpec : JSON.stringify(version.workflowSpec)) : {} const configuration = version.configuration ? JSON.parse(JSON.stringify(version.configuration)) : null return { graph, spec, configuration, version } } function parseMaybeJsonObject(value: T | string | null | undefined, fallback: T): T { if (value === null || value === undefined) return fallback if (typeof value === 'string') return JSON.parse(value) as T return value } function hasPayloadField( payload: UpdateWorkflowVersionPreservingStatePayload, key: K, ): boolean { return Object.prototype.hasOwnProperty.call(payload, key) } export async function updateWorkflowVersionPreservingState( versionId: string, payload: UpdateWorkflowVersionPreservingStatePayload, ): Promise { const existing = await getPromptVersion(versionId) const workflowSpec = hasPayloadField(payload, 'workflowSpec') ? payload.workflowSpec : parseMaybeJsonObject(existing.workflowSpec, {}) const graph = hasPayloadField(payload, 'graph') ? payload.graph : parseMaybeJsonObject(existing.graph, { nodes: [], edges: [] }) const variables = hasPayloadField(payload, 'variables') ? payload.variables : hasPayloadField(payload, 'workflowSpec') ? undefined : existing.variables return await updateWorkflowVersion(versionId, { title: hasPayloadField(payload, 'title') ? payload.title : existing.title, workflowSpec, graph, variables, configuration: hasPayloadField(payload, 'configuration') ? payload.configuration : existing.configuration, content: hasPayloadField(payload, 'content') ? payload.content : existing.content, markdownContent: hasPayloadField(payload, 'markdownContent') ? payload.markdownContent : existing.markdownContent, promoted: hasPayloadField(payload, 'promoted') ? payload.promoted : existing.promoted, draft: hasPayloadField(payload, 'draft') ? payload.draft : existing.draft, version: payload.version, strictUiCompatibility: payload.strictUiCompatibility, strictReferenceCompatibility: payload.strictReferenceCompatibility, strictAuthoring: payload.strictAuthoring, }) } /** * Patch a single node in an existing workflow version by name. * Creates a new version with the patched node — does not modify the original. * * Usage: * await patchWorkflowNode({ * promptId: 'prompt-uuid', * versionId: 'version-uuid', * nodeName: 'Fetch GitHub PRs', * patch: { code: 'async function execute(input) { ... }' }, * title: 'v3 — fix regex', * }) */ export async function patchWorkflowNode(params: { promptId: string versionId: string nodeName: string patch: Record title?: string version?: WorkflowVersion strictUiCompatibility?: boolean strictReferenceCompatibility?: boolean strictAuthoring?: boolean }): Promise { const { graph, spec, configuration, version } = await forkVersion(params.versionId) const workflowVersion = normalizeVersion(params.version) const { updateWorkflowNodeInput } = await import('./workflow-updates.ts') const updatedGraph = updateWorkflowNodeInput(graph, params.nodeName, params.patch, { version: workflowVersion, strictUiCompatibility: params.strictUiCompatibility, strictReferenceCompatibility: params.strictReferenceCompatibility, strictAuthoring: params.strictAuthoring, }) return await createWorkflowVersion({ promptId: params.promptId, title: params.title ?? `Patch: ${params.nodeName}`, workflowSpec: spec, graph: updatedGraph, variables: version.variables, configuration, content: version.content, markdownContent: version.markdownContent, version: workflowVersion, strictUiCompatibility: params.strictUiCompatibility, strictReferenceCompatibility: params.strictReferenceCompatibility, strictAuthoring: params.strictAuthoring, }) }