Chuyển đến nội dung chính

レッスン 13: ワークフローの検証と実行

DAG 検証: サイクル検出、孤立ノード。トポロジカルソートの実行。エラー処理、再試行、タイムアウト。実行履歴、ステップバイステップのデバッグ。 EventBus を介したリアルタイムの進行状況。

🧠 AI と ML — レッスン 12 レッスン 13: ワークフローの検証と実行

AIエージェントプラットフォームをゼロから構築 — xClawとの実戦

パート 4: ワークフロー エンジン

xdev.asia

はじめに

無効なワークフロー (サイクル、孤立ノードを含む) はランタイムをクラッシュさせます。この記事では、運用グレードのワークフロー エンジンの検証 + トポロジカル ソートの実行 + エラー処理を実装します。


1. ワークフローの検証

// packages/core/src/workflow/workflow-engine.ts
validateWorkflow(workflow: Workflow): ValidationResult {
  const errors: string[] = [];
  const warnings: string[] = [];

  // 1. Must have exactly one trigger node
  const triggers = workflow.nodes.filter(n => n.type === 'trigger');
  if (triggers.length === 0) errors.push('Workflow must have a trigger node');
  if (triggers.length > 1) errors.push('Workflow must have exactly one trigger node');

  // 2. Must have at least one end node
  const ends = workflow.nodes.filter(n => n.type === 'end');
  if (ends.length === 0) warnings.push('No end node — workflow may not terminate cleanly');

  // 3. Cycle detection (DFS)
  if (this.hasCycle(workflow)) {
    errors.push('Workflow contains a cycle — not allowed in DAG');
  }

  // 4. Orphan detection
  const connectedIds = new Set<string>();
  for (const edge of workflow.edges) {
    connectedIds.add(edge.source);
    connectedIds.add(edge.target);
  }
  const orphans = workflow.nodes.filter(n =>
    n.type !== 'trigger' && !connectedIds.has(n.id)
  );
  if (orphans.length > 0) {
    warnings.push(`Orphan nodes: ${orphans.map(n => n.label).join(', ')}`);
  }

  // 5. All edges reference valid nodes
  const nodeIds = new Set(workflow.nodes.map(n => n.id));
  for (const edge of workflow.edges) {
    if (!nodeIds.has(edge.source)) errors.push(`Edge references missing source: ${edge.source}`);
    if (!nodeIds.has(edge.target)) errors.push(`Edge references missing target: ${edge.target}`);
  }

  return {
    valid: errors.length === 0,
    errors,
    warnings,
  };
}

private hasCycle(workflow: Workflow): boolean {
  const adj = new Map<string, string[]>();
  for (const edge of workflow.edges) {
    if (!adj.has(edge.source)) adj.set(edge.source, []);
    adj.get(edge.source)!.push(edge.target);
  }

  const visited = new Set<string>();
  const inStack = new Set<string>();

  function dfs(nodeId: string): boolean {
    visited.add(nodeId);
    inStack.add(nodeId);

    for (const neighbor of adj.get(nodeId) || []) {
      if (inStack.has(neighbor)) return true; // Cycle!
      if (!visited.has(neighbor) && dfs(neighbor)) return true;
    }

    inStack.delete(nodeId);
    return false;
  }

  for (const node of workflow.nodes) {
    if (!visited.has(node.id) && dfs(node.id)) return true;
  }
  return false;
}

2. トポロジカルソートの実行

async execute(
  workflow: Workflow,
  triggerData: unknown,
  context: ToolContext,
): Promise<WorkflowResult> {
  // Validate first
  const validation = this.validateWorkflow(workflow);
  if (!validation.valid) {
    throw new Error(`Invalid workflow: ${validation.errors.join('; ')}`);
  }

  const ctx: WorkflowContext = {
    workflow,
    variables: { ...workflow.variables },
    triggerData,
    toolContext: context,
    nodeResults: new Map(),
    executionLog: [],
  };

  // Topological sort
  const executionOrder = this.topologicalSort(workflow);

  for (const nodeId of executionOrder) {
    const node = workflow.nodes.find(n => n.id === nodeId)!;
    const handler = this.handlers.get(node.type);

    if (!handler) {
      throw new Error(`No handler for node type: ${node.type}`);
    }

    const stepStart = Date.now();
    try {
      // Handle condition branching
      if (node.type === 'condition') {
        const result = await handler(node, ctx);
        ctx.nodeResults.set(nodeId, result);

        // Skip nodes not on the chosen branch
        // (handled by edge conditions)
        continue;
      }

      const result = await handler(node, ctx);
      ctx.nodeResults.set(nodeId, result);
      ctx.variables[`_node_${nodeId}`] = result;

      ctx.executionLog.push({
        nodeId,
        nodeType: node.type,
        status: 'success',
        result,
        duration: Date.now() - stepStart,
      });
    } catch (error) {
      ctx.executionLog.push({
        nodeId,
        nodeType: node.type,
        status: 'error',
        error: error instanceof Error ? error.message : String(error),
        duration: Date.now() - stepStart,
      });

      // Stop execution on error (or retry if configured)
      if (node.config.onError !== 'continue') {
        throw error;
      }
    }
  }

  return {
    success: true,
    variables: ctx.variables,
    executionLog: ctx.executionLog,
  };
}

private topologicalSort(workflow: Workflow): string[] {
  const inDegree = new Map<string, number>();
  const adj = new Map<string, string[]>();

  for (const node of workflow.nodes) {
    inDegree.set(node.id, 0);
    adj.set(node.id, []);
  }

  for (const edge of workflow.edges) {
    adj.get(edge.source)!.push(edge.target);
    inDegree.set(edge.target, (inDegree.get(edge.target) || 0) + 1);
  }

  const queue: string[] = [];
  for (const [id, degree] of inDegree) {
    if (degree === 0) queue.push(id);
  }

  const order: string[] = [];
  while (queue.length > 0) {
    const nodeId = queue.shift()!;
    order.push(nodeId);

    for (const neighbor of adj.get(nodeId) || []) {
      const newDegree = inDegree.get(neighbor)! - 1;
      inDegree.set(neighbor, newDegree);
      if (newDegree === 0) queue.push(neighbor);
    }
  }

  return order;
}

3. ワークフローの例: コンテンツ生成パイプライン

{
  "nodes": [
    { "id": "1", "type": "trigger", "label": "Manual Trigger", "config": {} },
    { "id": "2", "type": "llm_call", "label": "Generate Outline", "config": {
      "prompt": "Create an outline for: {{topic}}", "systemPrompt": "You are a content strategist"
    }},
    { "id": "3", "type": "llm_call", "label": "Write Article", "config": {
      "prompt": "Write article from outline: {{_node_2}}"
    }},
    { "id": "4", "type": "llm_call", "label": "Generate SEO", "config": {
      "prompt": "Generate SEO metadata for: {{_node_3}}"
    }},
    { "id": "5", "type": "end", "label": "Done", "config": {} }
  ],
  "edges": [
    { "source": "1", "target": "2" },
    { "source": "2", "target": "3" },
    { "source": "3", "target": "4" },
    { "source": "4", "target": "5" }
  ],
  "variables": { "topic": "AI Agent Architecture" }
}

4. まとめ

  • 検証 — サイクル検出、孤立ノード、欠落した参照
  • トポロジカルソート — 依存関係の順序に従ってノードを実行します
  • エラー処理 — ノードごとのエラー ポリシー (停止、続行、再試行)
  • 実行ログ — 段階的なデバッグと監査証跡

次の記事: スキル システム — 動的なツール構成。