ADK を使用したエージェント ワークフロー

1. はじめに

VibeStudio

この Codelab では、Agent Development Kit(ADK)のワークフローとグラフを使用して次世代のエージェント システムを構築する方法を説明します。一般的なアーキテクチャ パターンを実装し、人間参加型(HITL)のインタラクションをオーケストレートし、長時間実行される非同期実行を処理します。また、企業ナレッジベースと永続メモリを統合して、エージェントの動作をカスタマイズして進化させます。最後に、これらの機能を接続して、動画の自動生成パイプラインを推進します。

シナリオ

VibeTube でデジタル チャンネルを運営しており、アクティブな視聴者がいて、クリエイティブなアイデアのバックログが拡大している。各動画の制作には、トレンドの形式の調査、視聴者のフィードバックの統合、スクリプトの作成、ポリシー準拠の確認、動画クリップの生成など、複数の段階にわたる継続的な実行が必要です。生成モデルは個々のアセットのドラフトを作成できますが、一貫性のあるリリースを実現するには、オーケストレーションされたエージェント アーキテクチャが必要です。

このライフサイクルを自動化するために、VibeStudio を構築します。このエージェント パイプラインは、ルーティン調査を並行して実行し、人間参加型の承認用に厳選されたオプションを提示し、動画を生成する前に自動化されたポリシーゲートを適用し、本番環境の実行全体でコンテキストを保持します。

アイデアから公開されたクリップまで、構築するワークフロー

学習内容

10-summary

  • グラフ エンジニアリングの基盤: マルチステップ エージェント アーキテクチャには、明示的な制御フローと構造化された実行パスが必要です。エッジタプル、START エントリ ポイント、並列ファンアウト集約用の JoinNode、状態に基づいて実行を制御する決定論的ルーターノードを使用して、ADK Workflow を構築します。
  • エージェント モードとライフサイクル コールバック: 特殊なタスクには、明確な運用動作と決定的なガードレールが必要です。ADK Agent インスタンスは、chatsingle_turn、ツール対応の task モードをワークフロー ノードとして使用して構成し、before_model_callbackafter_agent_callback でインターセプタを適用します。
  • 人間参加型オーケストレーション: 重要なクリエイティブ チェックポイントで、人間が判断できるようにプロダクション パイプラインを一時停止します。RequestInput を実装して、ワークフローの実行を一時停止し、構造化されたレスポンス スキーマを適用し、アイドル状態のランタイム プロセスを維持せずに実行を再開します。
  • 階層型エージェント メモリ: 本番環境システムでは、一時的な実行状態と永続的なコンテキストが分離されます。Event(state=...) とパラメータ バインディングを使用して短期セッションの状態を管理し、GEAP メモリバンクを接続して、実行間でクリエイターの設定を抽出、統合、永続化します。
  • 企業のナレッジベースによるグラウンディング: 自律型エージェントには、動的なドメイン コンテキストとユーザーの感情が必要です。GEAP RAG Engine コーパスを並列ファンアウト内の専用の検索ノードとして接続し、エージェントの出力をセマンティックにグラウンディングします。
  • 長時間実行されるワークフローとデプロイ: マルチモーダル動画レンダリングは、長時間にわたって非同期で動作します。保留中の通話レシートを使用して LongRunningFunctionTool を実装し、通話 ID でワークフローを一時停止して再開します。また、ADK Runner を使用して、完成したパイプラインを Cloud Run にデプロイします。

この Codelab の構成

この Codelab は、コンセプトとアーキテクチャのリファレンスとして使用できます。各セクションでは、対応するワークベンチ ステップで実装される ADK 構成について説明し、参照コードを提供して、コア設計原則を確立します。各セクションを確認してから、ワークベンチで対応する演習を完了します。

ハンズオンは、インタラクティブなコードエディタ、ランタイム検証ツール、組み込みの ADK インスペクタを備えたコンパニオン ウェブ インターフェースである VibeStudio Workbench で行われます。ワークベンチのステップ番号は、この Codelab と直接対応しており、進捗状況を同期できます。基盤となるグラフの編集は手順間で保持され、ワークベンチでは手順を進めるにつれて前提条件が自動的に検証されます。

ワークベンチの演習を完了すると、エンドツーエンドのエージェント パイプラインを組み立て、実行中の VibeStudio アプリケーションを Cloud Run にデプロイして動画コンテンツを生成します。

実行場所: VibeStudio Workbench、バックエンド、Google Cloud サービス

この環境は、VibeStudio Workbench(コード編集とランタイム検証用のローカル ウェブ インターフェース)、バックエンド(ADK Workflowagent/ のステージ サンドボックス)、Google Cloud(Gemini モデル、GEAP メモリバンク、RAG Engine、Veo 動画生成)の 3 つの主要コンポーネントで構成されています。

2. セットアップ

ワークショップ クレジットを利用する

講師主導のラボに参加している場合は、講師が Google Cloud プロジェクトのクレジットを配布します。講師の指示に沿ってクレジットを適用し、アカウントで請求が有効になっていることを確認してから続行してください。

Cloud Shell を開く

Cloud Shell は、gcloud、Python、git がプリインストールされたブラウザベースの開発環境です。

Cloud Shell を起動するには:

  1. Google Cloud コンソールに移動します。
  2. 上部のナビゲーション ヘッダーで、[Cloud Shell をアクティブにする](ターミナル ウィンドウ アイコン)をクリックします。

Cloud Shell

ブラウザ ウィンドウの下部にターミナル セッションが開きます。

リポジトリのクローンを作成して初期化する

Cloud Shell ターミナルで次のコマンドを実行して、プロジェクトのクローンを作成します。

git clone https://github.com/gca-americas/vibetube-studio
cd ~/vibetube-studio

構成プロンプト

セットアップ中に、次の情報の入力を求められます。

  • Google Cloud プロジェクト ID: setup_project.sh から入力を求められたら、Enter キーを押して、新しいプロジェクトを自動的に作成します。既存のプロジェクト(事前に割り当てられたプロジェクトなど)を使用する場合は、プロジェクト ID を入力し、スペルが正しく、課金が有効になっていることを確認します。
  • イベントコード: 教師から提供された会議室コードを入力します。届いていない場合は、ティーチング アシスタントまたは近所の方に確認してください。このラボを自宅で完了する場合は、Enter キーを押してデフォルトの sandbox ルームを受け入れます。
  • チャンネルの表示名: setup_codelab.sh から求められたら、名前または希望するチャンネル ハンドルを入力するか、Enter キーを押して Google アカウントから生成されたデフォルトを受け入れます。

次の 2 つの設定スクリプトを順番に実行します。

./setup_project.sh
./setup_codelab.sh
  • setup_project.sh: 課金が有効な Google Cloud プロジェクトを作成または再利用し、プロジェクト ID を ~/project_id.txt に保存して、アクティブな gcloud コンテキストを構成します。
  • setup_codelab.sh: uv と Python の依存関係を .venv にインストールし、必要な Google Cloud API を有効にして、.env でチャネル設定を構成し、Gemini でモデル アクセスを確認し、メモリバンクと RAG リソースをプロビジョニングし、ワークベンチ インターフェースを構築して、VibeStudio ワークベンチを起動します。

スクリプトはプリフライト チェックを実行し、バックグラウンドで VibeStudio Workbench を起動します。最後の行に、開くためのリンクが表示されます。

7 · Preflight
   python 3.12
   auth path A: Vertex via ADC (STUDIO_VERTEX=1)
   Google Cloud ADC (project <your-project>)
   stage0_prompt loads
  ...
   stage6_video loads (13 edges)
   aiplatform.googleapis.com enabled (Gemini, Veo, Memory Bank, RAG Engine)
   vectorsearch.googleapis.com enabled (the vector store a RAG corpus is built on)
   Memory Bank connected
   RAG corpus connected
   VibeStudio Workbench running on port 4600

PREFLIGHT GREEN

Setup finished. The VibeStudio Workbench is already running.

  Open this and start at step 1
      https://4600-<your cloud shell host>/step/story

  It runs in the background. You do not need to start anything else.
      log      runs/lab.log
      stop     kill $(cat runs/lab.pid)
      start    scripts/start.sh

そのリンクをクリックします。同じアドレスは、[ウェブでプレビュー] > [ポートを変更] > [4600] でも確認できます。

環境を再確認するには、いつでも python scripts/preflight.py を実行します。ワークベンチを再起動するには、scripts/restart.sh を実行します。再度設定するには、./setup_codelab.sh を実行します。これにより、構成と進行状況が保持されます。

開いた状態で、シナリオのステップ 1: ストーリーと、完成したグラフの形状のステップ 2: 作成する内容を読みます。どちらにもエクササイズはありません。その後、このページに戻って手順 3 に進みます。

VibeStudio Workbench のハンズオン パートはすべて、実際のアーティファクト(ディスク上のファイルと実行によって書き込まれたセッション)を読み取る検証パネルで終わります。

リポジトリのレイアウト

リポジトリは、コア ワークフロー ロジック、ステップバイステップのサンドボックス、ワークベンチ環境、本番環境アプリケーションに構造化されています。

vibe-studio-lab/
├── agent/                  # Core ADK workflow, graph definition, and platform services
   ├── graph.py            # Workflow graph definition, node functions, and routers
   ├── desk.py             # Video render desk using LongRunningFunctionTool
   ├── schemas.py          # Pydantic schemas for directions, gates, and scripts
   ├── trends.py           # Trend generation and sampling utilities
   ├── backlog.txt         # Creator video ideas backlog
   ├── comments.md         # Audience comments for RAG Engine corpus seeding
   ├── policy_words.txt    # Blocked subject words for deterministic policy checks
   └── platform/           # Google Cloud service clients (Memory Bank, RAG, Veo)
       ├── config.py       # Environment variables, locations, and model configurations
       ├── memory.py       # GEAP Memory Bank callbacks and context injection
       ├── rag.py          # GEAP RAG Engine corpus creation and semantic retrieval
       └── videogen.py     # Veo video generation and operation polling
├── stage0_prompt/          # Step sandboxes: isolated agent.py files runnable in adk web
   └── ...                 # stage1_fanout through stage6_video for incremental steps
├── server/ & web/          # VibeStudio Workbench (FastAPI backend and React frontend)
├── vibestudio/             # Complete production application deployed to Cloud Run
   ├── server/             # FastAPI production server and event runner
   ├── web/                # End-user React web application
  • agent/: コア ワークフロー グラフが含まれます。このディレクトリ内のファイルを編集して、並列ファンアウト ノード、決定論的ポリシー ルーティング、メモリ コールバック、動画生成ツールを実装します。
  • agent/platform/: Gemini モデル、GEAP メモリバンク、GEAP RAG Engine、Veo 動画合成などの Google Cloud サービスと連携します。
  • stage0_prompt/stage6_video/: 自己完結型のサンドボックス環境。各フォルダはスタンドアロンの root_agent をエクスポートするため、埋め込み ADK 開発インターフェースを介して各ステップを個別に実行して検査できます。
  • server/web/: ポート 4600 でローカルに実行されている VibeStudio Workbench アプリケーション。ステップのドキュメント、ページ埋め込みコードエディタ、ランタイム証拠検証ツール、グラフの可視化をホストします。
  • vibestudio/: 最終ステップでパッケージ化され、Cloud Run にデプロイされる完全な本番環境アプリケーション。完了したワークフロー グラフのスタンドアロン コピーが含まれています。

3. モノリシック エージェント

マルチノード ワークフロー グラフを構築する前に、stage0_prompt/agent.py の単一のエージェントを使用してアーキテクチャのベースラインを確立します。このエージェントは、2 つの Python 関数ツールでサポートされている、本番環境パイプラインを文章で説明するモノリシック システム プロンプトに依存しています。

このベースラインを評価することで、プロンプト駆動型調整の運用上の境界が明らかになり、本番環境システムでグラフ オーケストレーションが必要な理由が明確になります。

ADK エージェントのアーキテクチャ(3A)

VibeStudio Workbench で、ステップ 3 · モノリシック エージェントに移動し、ADK エージェント アーキテクチャ(3A)を開きます。このビューには、ADK エージェント(LlmAgent)のコア アーキテクチャ レイヤが表示されます。

03-3A

from google.adk.agents import LlmAgent
from google.adk.tools import mcp_toolset

root_agent = LlmAgent(
    model="gemini-3.5-flash",                 # model
    instruction=BRAND_INSTRUCTION,            # instruction
    skills=[load_skill("brand-audit")],       # skills
    tools=[mcp_toolset("mcp_brand_style")],   # tools
    output_schema=BrandStyleReport,           # structured output
    before_agent_callback=setup_ctx,          # interceptor
    before_model_callback=require_image,      # interceptor
    after_model_callback=schema_guard,        # interceptor
)

インタラクティブな図では、エージェント コンポーネントが 5 つの運用ドメインにグループ化されています。

  • 推論レイヤ(モデル): 認知タスク、プロンプト推論、ツール選択を実行するコア言語モデル(Gemini 3 Flash など)。アーキテクチャの他のすべての要素は、このモデルに情報を提供するか、このモデルを制約します。
  • コンテキスト レイヤ(指示とスキル): モデルの推論を形成するディレクティブ。instruction は、永続的なシステム プロンプト、ペルソナ、運用ルールを確立します。skills は、繰り返し可能なワークフローのバージョン管理された手順ガイダンス(SKILL.md)を提供します。
  • コラボレーションとアクション レイヤ(ツール、サブエージェント、ワークフロー、出力スキーマ): エージェントが外部システムでアクションを実行し、型付きデータを出力できるようにするインターフェース。tools は、呼び出し可能な Python 関数または Model Context Protocol(MCP)エンドポイントを提供します。subagents は、下位の委任されたタスクを実行します。workflow は、マルチエージェント グラフを調整します。output_schema は、下流のコンシューマーが非構造化テキストではなく検証済みの JSON を受け取ることを保証するために、Pydantic モデルを適用します。
  • インターセプタ レイヤ(ライフサイクル コールバック): エージェントの実行(before_agent/after_agent)、個々のモデルのターン(before_model/after_model)、ツール呼び出し(before_tool/after_tool)の前後にカスタムコードを実行する決定論的なガードレール。インターセプタは、モデルのコンプライアンスに依存せずにポリシー ルールを適用します。
  • 外部状態(セッションとメモリ): エージェント ロジックから分離されたステートフルな永続性。Session は、現在の実行スレッドの一時的なワーキング メモリとイベント トレースを保持します。Memory は、GEAP メモリバンクなどのマネージド サービスを使用して、セッション間で永続的な事実と設定を維持します。

このステップのモノリシック エージェントは、これらのプリミティブのうち 3 つ(modelinstructiontools)のみを実装します。以降の手順では、グラフ ワークフロー、構造化スキーマ、インターセプタ、永続メモリ サービスについて説明します。

モノリシック エージェントの仕様(3B)

ワークベンチで、[モノリシック エージェントの仕様(3B)] に進みます。stage0_prompt/agent.py を開いて、ベースライン エージェントの定義を確認します。

  • 単一のプロンプト指示: システム プロンプトは、プラットフォームのトレンドの発見、バックログのアイデアの確認、クリエイティブ コンセプトの提案、禁止されているテーマに関するポリシーの適用、ショットリストの作成という 5 つの異なる制作タスクを連続した文章に凝縮します。
  • 基盤となるデータソース: エージェントは、グラフの横に定義された 2 つのソースを参照します。
    • agent/trends.py: 250 個のプールから、動的なヒートスコアを使用して 10 個のアクティブな形式とスタイルのトレンドをサンプリングします。
    • agent/backlog.txt: クリエイターのコンセプトのメモの元データを 1 行ずつ読み取ります。

エージェントのツール(3C)

ワークベンチで、[Tools in Agent (3C)](エージェントのツール(3C))に進みます。

エージェントにとってのツールとは

言語モデルは本質的に閉じた世界の推論エンジンです。事前にトレーニングされた重みと、そのコンテキスト ウィンドウに存在するトークンのみに基づいて動作します。データベースのクエリ、リアルタイム API へのアクセス、コードの実行をネイティブに行うことはできません。

ツールはこの境界を橋渡しします。これにより、モデルに外部エージェンシーが付与され、外部システムでグラウンド トゥルース情報を取得して決定論的アクションを実行できるようになります。

03-3C

ツール呼び出しは、モデルと ADK ランタイム間の明示的な 5 段階のプロトコルに従います。

  1. スキーマ宣言: デベロッパーが Python 関数をエージェントに提供します。ADK は、各関数の名前、型アノテーション、docstring を検査して、パラメータとその目的を記述する OpenAPI 互換の JSON スキーマ宣言を生成します。
  2. モデルの推論: 推論中に、モデルはユーザーのプロンプトに外部データが必要かどうかを評価します。必要に応じて、モデルはスキーマに一致するターゲット関数名と引数ディクショナリを含む構造化された function_call イベントを出力します。
  3. ランタイム実行: モデル自体はコードを実行しません。ADK ランタイムは function_call をインターセプトし、指定された引数を使用して実際のローカル Python 関数を実行し、戻り値を取得します。
  4. コンテキストの再注入: ADK ランタイムは、関数の戻り値を function_response イベントにパッケージ化し、アクティブなセッション履歴に追加します。
  5. 最終的な合成: モデルは、コンテキスト ウィンドウに存在するツールの出力を処理し、レスポンスを完了します。

stage0_prompt/agent.py では、2 つの調査ツールが標準の Python 関数として定義されています。

def check_trends() -> dict:
    """Ten formats trending on the platform right now, with a heat score each."""
    from agent.trends import sample_trends
    return {"trends": sample_trends()}


def read_backlog() -> dict:
    """The creator's backlog: ideas they noted down to make someday."""
    from agent.graph import backlog_notes
    return {"backlog": backlog_notes()}

ハンズオン編集と実行

ワークベンチのコードエディタで、エージェントの tools リストに 2 つの関数参照を追加します。

    tools=[check_trends, read_backlog],

変更を保存します。ディスク上のファイルが更新され、検証行で両方のツールが接続されていることが確認されます。

[Open adk web] をクリックして、埋め込み ADK 開発インターフェースを起動します。提案されたアイデアのプロンプトを送信します。

tonight's idea: a tiny robot doing laundry at midnight

想定される結果とその理由

このプロンプトを送信すると、セッション トレースで次の実行シーケンスが確認できます。

  • レスポンスの前に 2 つのツール実行イベントが表示される: check_trendsread_backlogfunction_call イベントと function_response イベントが表示されます。
    • 理由: Gemini はシステム プロンプト ディレクティブ(「トレンドを確認し、アイデアのバックログを確認する」)を評価し、重みにプラットフォームのトレンドとチャネルのメモがないことを認識し、両方の関数を呼び出してコンテキストをグラウンディングしました。
  • エージェントが方向性を提案し、確認を求める: トレンドとバックログを統合した動画の方向性が提案され、確認を求められます。
    • 理由: 指示ディレクティブで、スクリプトを生成する前にクリエイターと方向性を合意するようモデルに指示したため。
  • フォローアップ ターンで確認をバイパスする: 2 つ目のメッセージ skip the questions, just describe the video を送信します。エージェントはすぐに確認をバイパスし、タイトルとショットの下書きを作成します。
    • 理由: プロンプトの指示は、決定論的な障壁ではなく、推奨されるガイドラインです。モノリシック エージェントでは、外部ワークフローが実行フローを制御しないため、ユーザー指示が既存のシステム プロンプト ルールをオーバーライドできます。

モノリシック プロンプトのアーキテクチャ上の制限事項

単一のプロンプトで分離されたデモに許容可能な出力を生成できますが、ワークベンチ検証ツールで境界条件をテストすると、企業にとって重要な制限が明らかになります。

  • 非構造化された調査の集約: ツールの実行順序は非決定的です。モデルは取得したデータを自由形式の文章に要約するため、ダウンストリーム システムで特定の主張を行ったソースを特定することはできません。
  • 未確認のポリシーの適用: モデルが独自の安全性コンプライアンスを評価します。モデルがトピックを安全と判断した場合、外部の決定論的ロジックはその検出結果を検証しません。
  • 強制されない人間参加型の一時停止: クリエイターの確認を求めるプロンプトの指示は推奨事項です。モデルに質問をバイパスするよう指示するフォローアップ メッセージを送信すると、人間の承認が完全にスキップされます。

このようなアーキテクチャのギャップがあるため、次のステップで構築する明示的なグラフ ワークフローにモノリシック エージェントを分解する必要があります。

4. エージェント ワークフローの基本

VibeStudio Workbench で、ステップ 4 · エージェント ワークフローの基本のパート 4A4D に移動します。

このステップでは、単一エージェントのベースラインから、ADK Workflow を使用した決定論的グラフ オーケストレーションに移行します。並列研究ファンアウトを構築し、結合ノードでブランチを同期し、スキーマ検証済みのクリエイティブ候補を生成し、決定論的な人間参加型の承認ゲートを導入します。

グラフ アーキテクチャと実行チェーン(4A)

ワークベンチで、グラフ アーキテクチャと実行チェーン(4A)を開きます。

ADK Workflow は、エッジリストで定義された有向グラフとしてエージェントの実行を構造化します。

  • チェーン: 順次タプルは、線形ノード実行((node_a, node_b, node_c))を定義します。
  • 並列ブランチ: 共通の開始ノードを共有する独立したチェーンが同時に実行されます。
  • 同期: JoinNode に収束するチェーンは、すべての受信ブランチがレポートを送信するまで待機してからリリースします。
  • 決定論的制御: 実行フローは、プロンプト テキストから推測されるのではなく、宣言されたコード構造に則って制御されます。

04-4A

ADK のノード アーキタイプ

ADK ワークフローは、複数の特殊なノードタイプで構成されます。各アーキタイプは、グラフ内で特定の運用上の役割を果たし、決定論的コード実行と生成モデルの推論を分離します。

ノード アーキタイプ

実装

パイプラインでの役割

関数ノード

Event を返す Python 関数

決定論的ロジック、データ取得、状態の変更を実行します。

結合ノード

組み込みの JoinNode インスタンス

同時実行ブランチを集約された辞書に同期します。

エージェント ノード

single_turn モードで実行中の Agent

上流の入力に対して指示を評価し、検証済みのデータを出力します。

ルーターノード

route タグを含む Event を返す関数

条件付きロジックを評価して、ダウンストリームの実行ブランチを選択します。

手入力ノード

RequestInput を生成する関数

外部ユーザーからの応答が届くまで実行状態を一時停止します。

root_agent = Workflow(
    name="stage1_fanout",
    description="2 real readers -> join -> one research dict",
    edges=[...])

この構成では、root_agent はスタンドアロンの Agent ではなく、Workflow のインスタンスです。ADK はワークフローをファーストクラス エージェントとして扱い、グラフ全体を統一されたアプリケーションとして読み込み、提供、検査できます。name は ADK Web にアプリケーションを登録し、edges リストは実行トポロジを定義します。

並列調査のファンアウト(4B)

ワークベンチで、並列研究ファンアウト(4B)に進みます。stage1_fanout/agent.py を開きます。

04-4B

関数ノードと同期バリア

調査フェーズでは、agent/graph.py からインポートされた 2 つの関数ノードを使用します。

  • scan_trends: 10 個のスコア付きプラットフォーム トレンドを含む Event(output={"trends": [...]}) を返します。
  • read_backlog: 最初の実行プロンプトとともに、15 個のチャネル バックログ アイデアを含む Event(output={"backlog": [...], "idea": "..."}) を返します。

各関数は node_input(前のノードの出力)を受け取り、Event を返します。

JoinNode は同期バリアとして機能します。すべてのインバウンド チェーンがイベントを配信するまで一時停止し、すべてのブランチ結果をノード名({"scan_trends": {...}, "read_backlog": {...}})でキー設定された辞書に集約します。

ハンズオン編集: 結合エッジと平行エッジを定義する

stage1_fanout/agent.py で、JoinNode をインスタンス化し、START から始まる 2 つの並列チェーンを接続します。

join_research = JoinNode(name="join_research")
    edges=[(START, scan_trends, join_research),
           (START, read_backlog, join_research)])

変更内容を保存後、ワークベンチの検証ツールは、結合とエッジが配線されていることを確認します。[Run Stage 1] を使用するか、埋め込み ADK ウェブ インターフェースを使用してステージを実行します。

想定される結果とその理由

  • 同時読み取り実行: 実行グラフでは、scan_trendsread_backlog が同時に実行されます。
    • 理由: 両方のチェーンは START で始まります。ADK エンジンは、独立したブランチを同時にスケジュールします。
  • 集約された辞書出力: ワークフローは join_research で完了し、両方のリーダーのエントリを含む辞書を出力します。
    • 理由: JoinNode は、後続のノードの実行を許可する前に、データの完全なキャプチャを保証します。

エージェント ノード(4C)

ワークベンチで、[Agent nodes (4C)](エージェント ノード(4C))に進みます。stage2_direction/agent.py を開きます。

04-4C

動作モードと構造化スキーマ

Workflow 内に埋め込まれている場合、Agent はデフォルトで single_turn モードで実行されます。

  • 前のノードの出力をコンテキスト入力として受け取ります。
  • 会話のやり取りなしで、単一の推論呼び出しを実行します。
  • 構造化データを次のノードに出力します。

output_schema=Directions を割り当てることで、エージェントはモデル出力に対して Pydantic 検証を適用します。ダウンストリーム グラフは、非構造化された文章ではなく、型付きオブジェクトを受け取ります。

class Direction(BaseModel):
    title: str           # <=60 chars, filmable, characterful
    angle: str           # the twist, one line
    hook: str = ""       # 2-4 words, the video's sticker line
    evidence: list[Evidence]


class Directions(BaseModel):
    candidates: list[Direction]   # exactly 4

PROPOSE_INSTRUCTION は、傾向とバックログの両方の証拠を引用して、4 つの候補を提案するようにモデルに指示します。候補 1 ~ 3 は、実現可能なチャネル コンセプトを提供しています。候補 4 は、次のステップで安全ゲートをテストするために、意図的にポリシー違反のコンセプトを導入しています。

ハンズオン編集: エージェント ノードを定義して結合をチェーンする

stage2_direction/agent.py で、propose_directions を構成し、ワークフロー エッジを拡張します。

propose_directions = Agent(
    name="propose_directions",
    model=config.MODEL,
    instruction=PROPOSE_INSTRUCTION,
    output_schema=Directions)
    edges=[(START, scan_trends, join_research),
           (START, read_backlog, join_research),
           (join_research, propose_directions, direction_gate)])

想定される結果とその理由

  • 辞書の直接使用: propose_directions は、join_research によって出力された JSON ペイロードを手動でフォーマットすることなく使用します。
  • 型付き候補出力: エージェントは、4 つの個別の候補を含む検証済みの Directions オブジェクトを出力します。ダウンストリーム ノードは、文字列解析なしで属性名(candidate.title)でフィールドを読み取ります。

人間参加型(4D)

ワークベンチで、[人間参加型(4D)] に進みます。agent/graph.py を開きます。

04-4D

プロンプトの指示と確定的な停止

費用が発生する制作ワークフローやコンテンツの公開には、重要な意思決定ポイントで人間による監督が必要です。単一のプロンプトでは、確認リクエストはユーザーがモデルに簡単にバイパスを指示できるアドバイザリ指示です。ADK ワークフローでは、実行エンジンによって人間の承認が強制されます。グラフは指定されたノードで停止し、外部のスキーマ検証済み入力が受信されるまで進行できません。

  • RequestInput を生成すると、ワークフローの実行が直ちに一時停止します。
  • ADK は、セッション ストアにオープン割り込み呼び出しを記録し、一意の interrupt_id を発行します。
  • 実行プロセスは、トークンやサーバー スレッドを消費せずに停止します。
  • グラフの実行は、スキーマと割り込み ID に一致する有効な function_response が送信された場合にのみ再開されます。

ハンズオン編集: RequestInput で実行を一時停止する

agent/graph.py で、direction_gate 内に一時停止呼び出しを実装します。

    yield RequestInput(
        message="Pick tonight's direction: 1, 2, 3 or 4.",
        response_schema={
            "type": "object",
            "properties": {
                "pick": {"type": "string", "enum": ["1", "2", "3", "4"]}}},
        payload={"candidates": cands})

RequestInput は、次の 3 つの属性を構成します。

  • message: ユーザーに表示されるレビュー プロンプト。
  • response_schema: フロントエンドが入力フォームとしてレンダリングし、送信時に ADK によって検証される JSON スキーマ。
  • payload: リクエストにバンドルされたメタデータ(4 つの候補)。クライアント インターフェースは、セッションの状態をクエリせずにレビューカードをレンダリングできます。

想定される結果とその理由

  • ワークフローが direction_gate で停止する: ADK Web またはワークベンチ インターフェースで、実行が一時停止し、インタラクティブな候補選択フォームが表示されます。
    • 理由: エンジンが RequestInput を検出し、実行状態を runs/sessions.db に永続化しました。
  • 再開には構造化された入力が必要: 任意のチャット テキストを送信しても、グラフは進みません。オプション(1、2、3、4)を選択すると、response_schema を満たす入力された function_response が送信され、実行が再開されます。

5. 状態とルーター

VibeStudio Workbench で、ステップ 5 · 状態とルーター(5A) から (5C) までの部分に移動します。

ユーザーの選択をセッション状態に保持し、決定論的ルーターノードを使用してチャネルの安全に関するポリシーを適用し、動画スクリプトを生成する前にポリシー違反を自動的に修正する反復タスク エージェントを組み立てます。

ワークフローの状態(5A)

ワークベンチで、ワークフローの状態(5A)に移動します。

05-5A

セッション状態とノード出力

ADK ワークフローでは、データは 2 つの異なるメカニズムを介してグラフを移動します。

  • ノード出力(Event(output=...): エッジリストで定義された直近の下流コンシューマーに厳密に転送されるデータ。
  • セッション状態(Event(state=...): 実行ライフサイクルの後続のノードからアクセスできる共有 Key-Value 辞書。

05-5A

ユーザーが direction_gate で候補を選択すると、選択は数値インデックス({"pick": "2"})として届きます。下流ノードには、タイトル、ナラティブ アングル、フックラインなど、完全な方向オブジェクトが必要です。persist_direction は、すべての中間ノード ペイロードを介して詳細なメタデータを渡すのではなく、解決された候補を共有セッション状態に書き込みます。

ノードはセッション状態の辞書全体を渡す必要はありません。ノードが Event(state=...) を生成すると、新しい Key-Value ペアまたは更新された Key-Value ペアのみが提供されます。ADK は、これらの更新をセッション ストアに自動的にマージします。

    yield Event(state={"direction": chosen["title"], "angle": chosen.get("angle", ""),
                       "hook": hook, "user:prefs": {"last_direction": chosen["title"]}})

この Event を生成すると、制御が Workflow ランタイムに渡され、新しい値が runs/sessions.db のセッション ジャーナルに保持されます。

パラメータ バインディング

ADK 関数ノードは、パラメータ検査を通じてセッション状態を自動的に読み取ります。関数シグネチャで既存の状態キーと一致するパラメータ名が宣言されている場合、ADK は状態からそのキーを抽出し、直接渡します。

def persist_direction(node_input, candidates: list = []):
    ni = node_input if isinstance(node_input, dict) else {}
    raw = ni.get("pick")
    pick = str(raw).strip() if raw is not None else ""
    if candidates:
        i = int(pick) - 1 if pick.isdigit() else 0
        chosen = candidates[max(0, min(len(candidates) - 1, i))]
    else:
        chosen = {"title": "untitled", "angle": "", "evidence": []}
    hook = chosen.get("hook") or " ".join(chosen["title"].split()[:4])

ここで、candidatesdirection_gate によってセッション状態に書き込まれました。ADK は、明示的な辞書ルックアップを必要とせずに、persist_direction(node_input, candidates: list = []) に直接バインドします。

user: で始まるキーは、ユーザーレベルのストレージでセッション間で保持されるため、後続のワークフロー実行でクリエイターの設定にアクセスできます。

ハンズオン編集: 状態の永続化とノードの接続

  1. agent/graph.pypersist_direction 内で、TODO: PERSIST_STATE の行を状態イベントの yield に置き換えます。
    yield Event(state={"direction": chosen["title"], "angle": chosen.get("angle", ""),
                       "hook": hook, "user:prefs": {"last_direction": chosen["title"]}})
  1. stage3_router/agent.py で、edges リストの 3 番目のチェーンに persist_direction を追加します。
           (join_research, propose_directions, direction_gate,
            persist_direction)

ファイルを保存します。ワークベンチで、state write in placepersist_direction in the chain の両方に緑色のチェックマークが表示されていることを確認します。

ルーターノード(5B)

ワークベンチで、ルーターノード(5B)に移動します。

05-5B

決定論的ポリシー ルーティング

ルーターは、上流の出力を評価し、条件付きグラフのブランチに沿って実行を指示する特殊な関数ノードです。生成エージェントとは異なり、ルーターは LLM 呼び出しを行わずに決定論的ロジックを実行します。

ルーターは route タグを指定する Event を返します。

def length_check(node_input):
    too_long = len(node_input.get("title", "")) > 60
    return Event(output=node_input, route="TRIM" if too_long else "PASS")

ワークフロー定義では、辞書として定義されたエッジ ターゲットが、ルート名を宛先ノードにマッピングします。

    (length_check, {"TRIM": shorten, "PASS": scripter}),

ワークフロー ルーター policy_check は、agent/policy_words.txt から禁止フレーズを読み取り、選択した方向のタイトルと角度に対して単語全体のマッチングを実行します。

    return Event(output=node_input, route="BLOCK" if bad else "OK")

ポリシーをハードコードされた手順ではなくデータとして保存することで、ワークフロー グラフを変更せずに更新できます。テキスト ファイルを更新すると、後続の実行にすぐに適用されます。評価は決定論的な正規表現照合であるため、生成スクリプトが開始される前に、トークン費用なしでミリ秒単位で実行されます。

リンク先: Scripter と Quarantine

ルーターは、2 つの下流ノードのいずれかにトラフィックを転送します。

  • scripter: 承認された指示を Script Pydantic スキーマに準拠した構造化された制作スクリプトに変換する single_turn エージェント ノード:
scripter = Agent(
    name="scripter",
    model=config.MODEL,
    instruction=SCRIPT_INSTRUCTION,
    output_schema=Script)
  • quarantine: 最初は、フラグが設定されたルートを停止するプレースホルダ関数でしたが、次の部分で自律的な修復エージェントに置き換えられます。

ハンズオン編集: ポリシー チェックのルーティング

  1. agent/graph.pypolicy_check 内で、return 文を完成させます。
    return Event(output=node_input, route="BLOCK" if bad else "OK")
  1. stage3_router/agent.py で、edges を更新して policy_check をルーティングし、検疫ブランチを scripter に再結合します。
           (join_research, propose_directions, direction_gate,
            persist_direction, policy_check),
           (policy_check, {"OK": scripter, "BLOCK": quarantine}),
           (quarantine, scripter)])

ファイルを保存します。ワークベンチで、ルーター エッジ マッピングが検証されていることを確認します。

エージェント モードとタスクノード(5C)

ワークベンチで、[Agent modes and the task node (5C)](エージェント モードとタスクノード(5C))に移動します。

05-5C

エージェントの実行モード

ADK Agent インスタンスは、特定のパイプライン要件に合わせて調整された 3 つの実行モードをサポートしています。

モード

実行ライフサイクル

パイプラインでの役割

chat

マルチターンの会話ループ。モデルは、ツールを呼び出すタイミング、入力を求めるタイミング、ターンの終了タイミングを決定します。

インタラクティブなユーザーに面しているルート エージェント。

single_turn

単一モデルの推論呼び出し。前のノードの入力を受け取り、構造化されたスキーマ オブジェクトを出力します。

シーケンシャル グラフ変換(propose_directionsscripter)。

task

ツール実行による自律ループ。エージェントは、組み込みの finish_task ツールを呼び出すまで反復処理を行います。

複数ステップの修復と検査(quarantine)。

自律型ポリシーの修復

フラグが設定された方向を書き換えるには、修復の反復回数が可変であるため、task モードが必要です。エージェントはフラグ付きの指示を受け取り、find_policy_hits を呼び出して違反を検出し、suggest_replacement を介して承認済みの代替案をリクエストし、指示を書き換え、続行する前にクリーンさを検証します。

どちらのツールも、型付きシグネチャと docstring を使用して agent/cleanup_tools.py で定義されています。

def find_policy_hits(text: str) -> dict:
    """Which refused words appear in `text`. Matches whole words and phrases
    from agent/policy_words.txt, case-insensitive.

    Returns {"hits": [...], "clean": bool}. clean is true when hits is empty.
    """


def suggest_replacement(word: str) -> dict:
    """The channel's approved stand-in for a refused word, read from
    agent/policy_replacements.txt.

    Returns {"word", "replacement", "listed"}. When the word has no entry,
    listed is false and replacement is a hint to pick a gentle synonym.
    """

ハンズオン編集: 検疫タスク エージェントを組み立てる

stage3_router/agent.py で、プレースホルダの quarantine 関数をタスク エージェントの定義に置き換えます。

quarantine = Agent(
    name="quarantine",
    model=config.MODEL,
    instruction=QUARANTINE_INSTRUCTION,
    mode="task",
    tools=[find_policy_hits, suggest_replacement],
    output_schema=CleanedDirection,
)

タスクモードでは、エージェントにツールが提供され、finish_task を呼び出すことで実行が終了します。mode="task" が構成されている場合、ADK は finish_task を自動的に提供し、そのパラメータを output_schema から導出します。これにより、ノードはスクリプター ノードの入力スキーマに一致する型付きの CleanedDirection オブジェクトを生成します。

05-5C

想定される結果とその理由

ADK Web または VibeStudio Workbench で両方の実行パスをテストします。

  • 承認済みルート(候補 1、2、3):
    • 承認された候補ルートを選択すると、policy_check から scripterroute="OK")に直接ルーティングされます。
    • スクリプターは、Script スキーマに準拠した 3 ショットの制作スクリプトを生成します。
  • 検疫の修復ルート(候補 4):
    • 候補 4 には、フラグが付けられた語彙(「クリックベイト」、「バイラル ハック」)が含まれています。
    • policy_check ルートを quarantineroute="BLOCK")に移動します。
    • セッション トレースで、quarantinefind_policy_hits を呼び出し、違反ごとに suggest_replacement を呼び出し、タイトルを書き換え、finish_task を呼び出していることを確認します。
    • 実行は scripter に戻り、サニタイズされた指示からスクリプトを生成します。

6. メモリバンク

VibeStudio Workbench で、ステップ 6 · メモリバンクのパート (6A)(6B) に移動します。

現在のところ、ワークフローはセッション間でメモリを使用せずに動作します。各実行は最初から始まり、クリエイターが以前に選択した内容や好みのジャンルは認識されません。このステップでは、Vertex AI Agent Engine Memory Bank を接続して、実行間でクリエイターの好みを保存して取得します。

重要なのは、メモリがパイプライン ノードではなくエージェント ライフサイクル コールバックを介して統合されることです。メモリの抽出と取得は、中間データ ステージではなく個々のエージェントに提供されるため、コールバックをアタッチすると、クリーンで分離されたグラフ トポロジが維持されます。

メモリバンク(6A)

ワークベンチで、[メモリバンク (6A)] に移動します。

06-6A

管理対象のユーザーレベルのメモリ

メモリバンクは、長期的なユーザー メモリ用のマネージド サービスです。定義されたスコープ(ここではアプリケーション名とユーザー ID で識別)の下で、人物に関する事実を整理します。

SCOPE = {"app_name": config.APP, "user_id": config.USER}
TOPICS = {
    "CREATOR_TASTE": "Which video directions this creator picks and passes on, "
                     "and how that preference changes over time.",
    "CHANNEL_RULES": "Standing instructions the creator states for every video "
                     "(style, subjects to avoid, format rules).",
}

カスタム メモリトピックは、バンクが記録する内容の境界を定義します。

  • トピックの抽出: memories.generate 経由で新しい会話テキストが送信されると、サービスは各トピックの説明に対して抽出モデルを適用します。トピックと一致しないテキストは、思い出として表示されません。
  • 統合と重複除去: サービスは、新しく抽出された事実をエンベディングに変換し、スコープ内の既存のメモリと比較します。観測が既存のメモリと一致すると、サービスはそのメモリを更新します。新しい情報を表す場合、サービスは新しいエントリを作成します。この統合プロセスにより、トピックに関する複数のセッションが冗長なエントリを生成するのではなく、一貫性のある概要に統合されます。
  • 取得: ユーザー スコープで memories.retrieve を呼び出すと、保存されたファクトが古い順に返されます。

どちらのオペレーションも agent/platform/memory.py に実装されています。プロビジョニングされた銀行リソース名は、runs/memorybank.json にローカルでキャッシュに保存されます。

メモリバンクを設定する

ワークベンチ コントロールを使用するか、ターミナルで CLI コマンドを実行します。

  1. 銀行を接続してプロビジョニングします
    python -m agent.platform.bank
    
    Agent Engine インスタンスを作成し、CREATOR_TASTE トピックと CHANNEL_RULES トピックを構成します。
  2. Seed historical sessions(過去のセッションをシードする):
    python -m agent.platform.bank load
    
    4 つの歴史的なクリエイター セッション(スタイル制約のある 2 つの動物テーマ、1 つのガジェット テーマ、1 つの最近のファンタジー テーマ)を読み込みます。
  3. 統合された事実を検査する:
    python -m agent.platform.bank list
    
    出力を確認します。ナラティブ トランスクリプトが、構造化された事実の統合ステートメントに変換されていることに注目してください。

コールバック(6B)

ワークベンチで、[コールバック(6B)] に移動します。stage4_memory/agent.py を開きます。

06-6A

ADK エージェント ライフサイクル コールバック

コールバックは、Agent の引数として渡される関数です。ADK は、事前定義されたライフサイクル タイミングでコールバックを呼び出し、アクティブなコンテキストを渡します。None を返すと、通常の実行が継続されます。置換オブジェクトを返すと、オペレーションがオーバーライドまたはインターセプトされます。

06-6A

ADK には、次の 3 組のコールバックが用意されています。

コールバック ペア

呼び出しポイント

受信したパラメータ

戻り値の動作

before_agent_callback
after_agent_callback

エージェントのターン全体を囲む

CallbackContext(状態、セッション、呼び出し)

Content を返すと、エージェントの返信が置き換えられます。None は通常どおりに処理されます。

before_model_callback
after_model_callback

各 LLM 推論呼び出しを囲む

LlmRequest または LlmResponse

LlmResponse を返すと、モデル呼び出しがインターセプトまたはスキップされます。None を返すと、処理が続行されます。

before_tool_callback
after_tool_callback

各ツールの実行を囲む

ツールの定義、引数、結果

辞書を返すと、ツールの出力がオーバーライドされ、None が続行されます。

コールバックは、ワークフロー グラフに余分なノードを導入することなく、コンテキストの挿入、ガードレール、テレメトリー、キャッシュ ルックアップを行うためのクリーンな場所を提供します。

ハンズオン編集: 配線リコールとコールバックの記憶

  1. stage4_memory/agent.py で、propose_directions を更新して before_model_callback=recall_taste を関連付けます。
    output_schema=Directions,
    before_model_callback=recall_taste)

recall_taste は、Gemini が候補のルートを生成する直前に実行されます。メモリバンクからクリエイターの履歴を取得し、最も古い順に思い出をフォーマットして、送信 LlmRequest に追加します。このプロンプトは、チャンネル ルールを厳格な制約として扱いながら、候補 1 ~ 3 をクリエイターの現在の好みに近づけるようにモデルに指示します。

  1. stage4_memory/agent.py で、scripter を更新して after_agent_callback=remember_pick を関連付けます。
    output_schema=Script,
    after_agent_callback=remember_pick)

remember_pick は、scripter のターンが完了した後に実行されます。セッション状態から選択された方向を読み取り、クリエイターの決定を要約した簡潔なステートメントを合成し、memories.generate を呼び出してメモリバンクを更新します。

想定される結果とその理由

ワークベンチまたは ADK Web でコールバック拡張ワークフローをテストします。

  1. 空のプロンプトで実行を実行します。
    • セッション トレースで、LlmRequestpropose_directions を確認します。ファンタジー テーマと簡潔なペースを好むクリエイターの好みを詳しく説明するメモリ コンテキストが追加されています。
    • 提案された方向性を確認します。候補 1 ~ 3 は、トレンドで他のトピックが強調されている場合でも、クリエイターの過去の好みに沿っています。
  2. direction_gate で候補を選択します。
  3. scripter が完了したら、メモリバンクのレコードを確認します。
    python -m agent.platform.bank list
    
    銀行に最新の選択が反映され、以前の好みの記録と統合されます。

7. RAG Engine

VibeStudio Workbench で、ステップ 7 · RAG Engine のパート (7A)(7B) に移動します。

07-7A

公開された動画には、視聴者からのフィードバックが継続的に蓄積されます。agent/comments.md には 30 件の代表的なコメントが収集され、視聴者の称賛、スポンサー付きコンテンツのペースに関する批判、音声の好みなどが記録されます。このステップでは、Vertex AI RAG Engine を使用してこれらのコメントのインデックスを作成し、セマンティック検索をリサーチ ファンアウトに接続します。

ドキュメントの取得(7A)

ワークベンチで、[RAG Engine (7A)] に移動します。

メモリバンクと RAG Engine の比較

どちらのツールもワークフローを外部データに結び付けますが、アーキテクチャ上の目的は異なります。

ディメンション

メモリバンク

RAG Engine

主なユースケース

長期的なユーザー設定と運用ルール

大規模なドキュメント コレクションに対するセマンティック検索

スコープ

個々のユーザー ID とアプリケーション名にスコープ設定

すべてのユーザーの共有コーパス リソースにスコープ設定

データ処理

リアルタイムの抽出、エンベディング、セマンティック統合

ドキュメントのチャンク分割、ベクトル エンベディング、最近傍検索

グラフの統合

エージェント ライフサイクル コールバック(before_model_callbackafter_agent_callback

研究ファンアウトの専用関数ノード(read_feedback

07-7A

ドキュメントのチャンク化とエンベディング

RAG Engine は、テキストをセマンティック パッセージに分割し、そのベクトルをマネージド データベースに保存することで、ドキュメントのインデックスを作成します。

corpus = rag.create_corpus(
    display_name="vibestudio-feedback",
    description="Vibe Studio: what the audience wrote under the channel's past videos.",
    backend_config=rag.RagVectorDbConfig(
        rag_embedding_model_config=rag.RagEmbeddingModelConfig(
            vertex_prediction_endpoint=rag.VertexPredictionEndpoint(
                publisher_model="publishers/google/models/text-embedding-005"))))

rag.upload_file(
    corpus_name=corpus.name, path="agent/comments.md", display_name="comments.md",
    transformation_config=rag.TransformationConfig(
        chunking_config=rag.ChunkingConfig(chunk_size=120, chunk_overlap=20)))
  • チャンクサイズ: 120 トークンに設定され、20 トークンが重複しています。これにより、各パッセージから 2 ~ 3 個のコメントが取得され、各ベクトルが関連性のないフィードバック全体で意味が薄まることなく、まとまりのある感情を表すようになります。
  • エンベディング モデル: text-embedding-005 はテキストを高次元ベクトルに変換します。クエリが送信されると、モデルはクエリをベクトルに変換し、セマンティック距離に基づいて最も近い一致を見つけます。靴下を守る小さなドラゴンについてのコメントは、キーワードが完全に一致していなくても、魔法の生き物についてのプロンプトと一致します。

RAG コーパスの設定

ワークベンチのボタンまたはターミナル コマンドを使用して、コーパスを初期化します。

  1. コーパスを作成する:
    python -m agent.platform.rag
    
    マネージド ベクトル データベースをプロビジョニングし、リソース ID を runs/ragcorpus.json に記録します。
  2. コメントをアップロードしてインデックス登録する: チャンク構成で agent/comments.md をアップロードし、インデックス登録が完了するまで待機します。
  3. コーパスをクエリする: コメントと完全に一致する単語を含まないクエリを使用して類似性検索をテストします(たとえば、「小さな魔法の生き物」というクエリで、ドラゴンに関するコメントを取得します)。

検索ノード(7B)

ワークベンチで、[The third reader (7B)] に移動します。stage5_rag/agent.py を開きます。

07-7B

グラフノードとしての取得

オーディエンス フィードバックは、ワークフロー全体で共有される調査データを表します。個人のクリエイターのメモリとは異なり、視聴者の感情はトレンドやバックログ データとともに join_research に直接フィードされます。したがって、関数ノードとして実装されます。

07-7B

def read_feedback(node_input):
    """The third reader (step 7): what the audience wrote under past videos,
    the passages nearest to tonight's idea. Retrieval, not a model call."""
    from .platform import rag
    idea = idea_text(node_input)
    query = idea or "what viewers liked and what they complained about"
    try:
        hits = rag.retrieve(query)
    except Exception as e:
        print(f"  [rag] feedback unavailable ({str(e)[:80]})")
        return Event(output={"query": query, "feedback": [],
                             "note": "no corpus connected - run: python -m agent.platform.rag"})
    return Event(output={"query": query, "feedback": [h["text"] for h in hits]})

read_feedback は、ユーザーの最初のアイデアを抽出し、RAG Engine コーパスに対してベクトル クエリを実行します。取得したコメントを Event(output=...) ペイロードで出力します。

ハンズオン編集: 3 つ目のリーダーをファンアウトに配線する

stage5_rag/agent.py で、edges を更新して、join_research に入る 3 番目の並列ブランチとして read_feedback を追加します。

           (START, read_backlog, join_research),
           (START, read_feedback, join_research),

join_researchJoinNode であるため、すべての受信ブランチを同期し、scan_trendsread_backlogread_feedback がすべてイベントを送信するまで待機してから、集約されたバンドルをダウンストリームに渡します。

想定される結果とその理由

ワークベンチでワークフローを実行します。

  1. アイデアのプロンプト(「キッチンのカウンターを守るミニチュアのドラゴン」など)を送信します。
  2. 実行トレースで、3 つのリーダーノードがすべて同時に実行されていることを確認します。
  3. join_research を確認します。出力辞書に trendsbacklogfeedback が含まれています。
  4. propose_directions から生成された候補を調べます。モデルは視聴者のコメントを提案に取り込み、エビデンス フィールドで視聴者の感情を参照します。
  5. RAG 検索は決定論的(同じクエリに対して同じコメント パッセージが返される)ですが、生成提案ノードはクリエイティブなバリエーションを生成します。

8. Veo を使用した非同期動画生成

VibeStudio Workbench で、ステップ 8 · 動画のパート (8A)(8B) に移動します。

Google Veo で高画質の動画を生成するには、レンダリングごとに数分かかります。この期間中にグラフの実行をブロックすると、コンピューティング リソースが無駄になり、スレッド プールがロックされ、実行が HTTP 接続のドロップアウトにさらされます。このステップでは、ADK の LongRunningFunctionTool を使用して動画レンダリングを非同期にします。

長時間実行ツール(8A)

ワークベンチで、[A long-running tool (8A)] に移動します。stage6_video/agent.pyagent/deliver.py を開きます。

08-8A

同期ツールと長時間実行ツール

標準の ADK 関数ツールは、エージェント ターン内で同期的に実行されます。モデルがツールを呼び出し、戻りペイロードを待機して、結果を進行中のターンに組み込みます。

動画のレンダリングを 1 ターンで完了できない。代わりに、render_submit は生成ジョブを開始し、ステータス "pending" のオペレーション レシートをすぐに返します。

def render_submit(prompt: str) -> dict:
    """Submit one Veo render of `prompt`. Returns at once with a pending
    receipt; the clip is delivered later, to this call, by id."""
    receipt = videogen.start(f"{prompt} {videogen.NO_TEXT}")
    return {"status": "pending", "operation": receipt["operation"], "prompt": receipt["prompt"]}

LongRunningFunctionTool でラップすると、ADK は "pending" ステータスをインターセプトします。エージェントのターンが終了し、ワークフローがノードで一時停止し、保留中の通話メタデータ(通話 ID と領収書を含む)が runs/sessions.db に記録されます。実行プロセスは、アクティブなネットワーク接続やワーカー スレッドを維持せずに正常に終了します。

ハンズオン編集: レンダリング ツールをラップする

stage6_video/agent.py で、render_desk を更新して render_submitLongRunningFunctionTool でラップします。

    tools=[LongRunningFunctionTool(render_submit)])

通話 ID で再開する

ユニバーサル再開パターン

ADK は、ユーザーと外部ツールの両方に対して、ワークフローの一時停止と再開に同じメカニズムを適用します。

停止トリガー

Construct の開始

保存された一時停止状態

再開イベント

人間による判断

yield RequestInput(...)

セッション ストアで入力プロンプトを開く

FunctionResponse 停止通話 ID を含む

長時間実行ツール

LongRunningFunctionTool(...)pending を返す

セッション ストアでツール呼び出しを開く

FunctionResponse 停止通話 ID を含む

どちらのシナリオでも、ワークフローは完全に停止し、一致する FunctionResponse を含むイベントが外部ソース(ユーザー インターフェース、Webhook、バックグラウンド ワーカー)から到着した場合にのみ再開されます。

ハンズオン編集: 配送レスポンスを完了する

agent/deliver.py で、再開 FunctionResponse 部分を構築します。

    part = Part(function_response=FunctionResponse(
        id=row["call_id"], name=row["name"], response=response))

配信デーモンは、動画ファイルが生成されるまで Veo をポーリングし、この FunctionResponse をセッションにディスパッチします。ADK は通話 ID を照合し、次のノードでワークフローを直接再開します。完了したノードは再実行されず、エージェントは別の生成ターンを実行しません。

.envSTUDIO_REAL_VIDEO=0 を設定すると、モック レンダリングが有効になります。start はテスト領収書をすぐに返し、check は課金対象の Veo API 呼び出しを行わずに 5 秒で完了をシミュレートします。

パイプラインの統合(8B)

ワークベンチで、グラフの render_desk(8B)に移動します。stage6_video/agent.py を開きます。

パイプラインの終端ノードは store_video です。runs/state.json(配信プロセスが記録した場所)から完了したレンダリング情報を読み取り、動画の URL と生成ステータスを共有セッション状態に commit します。

08-8B

ハンズオン編集: 完全な動画パイプラインの配線

stage6_video/agent.py で、edges を更新して render_deskstore_video を追加します。

           (quarantine, scripter),
           (scripter, render_desk, store_video)])

想定される結果とその理由

ワークベンチで非同期生成フローをテストします。

  1. 候補の選択とスクリプトの生成を通じてワークフローを実行します。
  2. render_desk で、エージェントが render_submit を呼び出すことを確認します。
  3. ワークフローは直ちに一時停止します。ワークベンチまたは ADK Web で、保留中のステータスを確認します。セッションはオープン コール ID を保持しており、バックグラウンド プロセスはリソースを消費していません。
  4. ワークベンチ コンソールまたはターミナルで配信デーモンを実行します。
    python -m agent.deliver
    
    配信プロセスは、動画の準備が整うまで Veo をモニタリングし、準備が整ったら再開イベントをディスパッチします。
  5. ADK Web でセッションを更新します。実行が store_video で再開され、動画 URL がセッション状態にコミットされ、ワークフローが完了します。

9. Cloud Run にデプロイする

VibeStudio Workbench で、[ステップ 9 · デプロイ] に移動します。

専用のサンドボックスでパイプラインの各コンポーネントを開発して検証しました。このステップでは、完全な本番環境パイプラインを組み立てて Google Cloud Run にデプロイします。

09-9A

ADK ランナー

開発では、adk web がグラフをオーケストレートしました。本番環境では、アプリケーションは ADK の Runner クラスを使用してワークフローをホストします。

self._svc = DatabaseSessionService(db_url=config.DB_URL)
self._runner = Runner(app_name=config.APP, agent=wf, session_service=self._svc)

async for ev in self._runner.run_async(user_id=config.USER, session_id=run_id, new_message=message):
    self._absorb(ev)    # fold the ADK event into the run state, publish one app event

# the gate's answer and the render's delivery are the same call, with a function_response part
part = Part(function_response=FunctionResponse(id=call_id, name=name, response=response))
  • run_async: ワークフローの実行を駆動し、ノードの実行時にイベントを順次生成し、セッション サービスへの更新を永続化します。
  • 統合された再開: direction_gate でのユーザーの決定と、Veo からの動画配信の完了の両方が、run_async に送信された同一の FunctionResponse オブジェクトを介して実行を再開します。

本番環境のアプリケーション アーキテクチャ

vibestudio/ の本番環境アプリケーションには、完全なパイプラインが統合されています。

vibestudio/
  server/
    main.py                 FastAPI: application server, REST routes, static assets
    api.py                  REST API endpoints: run, pick, publish, backlog, profile, history
    runner.py               Runner orchestration over the workflow, background render poller
    platform/               Event bus (SSE stream), file storage, publishing, telemetry
    agent/                  Production agent package, verified by checks/verify_app.py
      graph.py              The complete workflow graph and node definitions
      desk.py               render_desk and render_submit wrapped with LongRunningFunctionTool
      schemas.py            Pydantic schemas: Directions, CleanedDirection, Script
      cleanup_tools.py      Deterministic policy tools: find_policy_hits, suggest_replacement
      platform/             Memory Bank, RAG Engine, and Veo integrations
  web/                      Production React user interface
  Dockerfile · deploy.py · run.sh
  • 単一のイベント ストリーム: FastAPI バックエンドは、単一のサーバー送信イベント(SSE)ストリーム全体でイベントをパブリッシュします。React フロントエンドは、グラフの進行状況をリアルタイムで可視化し、状態を失うことなく遅延接続を処理します。
  • 分離された実行: アプリケーションがイベントループを管理します。ワークフロー グラフは、実行ロジックに完全に焦点を当てており、フロントエンド インターフェースを認識していません。

agent/graph.py の完全なワークフロー エッジリストは、この Codelab で構築されたすべてのアーキテクチャ パターンを組み合わせたものです。

        (START, scan_trends, join_research),
        (START, read_backlog, join_research),
        (START, read_feedback, join_research),
        (join_research, propose_directions, direction_gate,
         persist_direction, policy_check),
        (policy_check, {"OK": scripter, "BLOCK": quarantine}),
        (quarantine, scripter),
        (scripter, render_desk, store_video),

Cloud Run へのデプロイ

Google Cloud Run は、自動スケーリング、リクエスト ルーティング、統合コンテナ ビルドを備えたサーバーレス ホスティングを提供します。

gcloud run deploy vibestudio --source vibestudio \
  --project $GOOGLE_CLOUD_PROJECT --region us-central1 \
  --labels dev-tutorial-codelab=vibetube --allow-unauthenticated \
  --memory 2Gi --cpu 2 --timeout 3600 --concurrency 40 \
  --max-instances 1 --min-instances 1 --session-affinity \
  --set-env-vars GOOGLE_CLOUD_PROJECT=...,STUDIO_VERTEX=1,STUDIO_MEMORY_BANK=...,STUDIO_RAG_CORPUS=...,VIBETUBE_URL=...,VIBETUBE_EVENT=...,VIBETUBE_NAME=...,VIBETUBE_PROJECT=...
  • コンテナ ビルド: gcloud run deploy --source は、vibestudio/ ディレクトリをパッケージ化し、Cloud Build を使用してコンテナ イメージをビルドし、1 回のオペレーションでサービスをデプロイします。
  • セッション アフィニティ: 同じユーザーからのリクエストを同じコンテナ インスタンスに転送し、反復ステップ間でローカル セッション状態を維持します。
  • オブザーバビリティ: Cloud Trace の統合により、すべてのノード、LLM 呼び出し、ツール実行の分散スパンが記録され、Google Cloud コンソールの Trace エクスプローラでアクセスできます。

ワークベンチの [デプロイ] ボタンをクリックして、デプロイ スクリプトを実行します。ビルドが完了すると、ターミナルにライブ サービス URL が表示されます。

アプリ

10. まとめ

VibeStudio Workbench で、ステップ 10 · 概要に移動して、完成したアーキテクチャを確認します。

10-summary

ステップ

アーキテクチャとコンセプト

実装パターン

単一のプロンプト

単一のプロンプト、関数ツール、順次チャットループ

Agent(tools=[...])function_call / function_response

エージェント ワークフローの基本

グラフ ワークフロー、並列調査、スキーマ出力、ヒューマン ゲート

WorkflowSTARTJoinNodeoutput_schemaRequestInput

状態とルーター

共有セッションの状態、パラメータ バインディング、決定論的ルーティング、タスク エージェント

Event(state=...)Event(route=...)mode="task"finish_task

メモリバンク

ユーザーレベルの長期記憶、セマンティック統合、ライフサイクル フック

memories.generate / retrievebefore_model_callbackafter_agent_callback

RAG Engine

視聴者のコメント、セマンティック エンベディングによるドキュメント取得

rag.create_corpusRagEmbeddingModelConfigread_feedback ノード

Veo を使用した非同期動画生成

長時間実行ツール、保留中の領収書、外部配信デーモン

LongRunningFunctionToolFunctionResponse(id=...) の再開

Cloud Run へのデプロイ

プログラムによるオーケストレーション、サーバー送信イベント、サーバーレス コンテナ

Runner(agent=wf)run_async、Cloud Run デプロイ

アーキテクチャの基本原則

  1. 待機ではなく一時停止: ワークフローは、人間の入力(RequestInput)または長時間実行オペレーション(LongRunningFunctionTool)のためにクリーンに一時停止します。プロセスは、スレッドまたはネットワーク ソケットでアイドル状態で待機しません。
  2. ユニバーサル再開: すべての一時停止は、一時停止されたノードの呼び出し ID を含む単一の function_response という同一のメカニズムで再開されます。
  3. 分離された状態管理: ノードは、冗長で密結合された中間ペイロードではなく、名前付きセッション状態キーとパラメータ バインディングを介してデータを共有します。
  4. 生成費用の前の決定論的ルーティング: ルールベースのルーターと正規表現フィルタは、生成モデルが実行される前に、ゼロトークン費用でポリシーを評価します。
  5. 関心の分離: 個々のエージェントに固有のコンテキストはライフサイクル コールバックに属し、共有データ依存関係は

10-output