コンテンツにスキップ

実行とライフサイクル (Execution and Lifecycle)

このページでは、ノードがワークフローのライフサイクルに参加するためにオーバーライドできるコールバックや、ポーリングを伴う外部 API 連携などの長時間実行される非同期処理のパターンについて解説します。

ライフサイクルコールバック (Lifecycle Callbacks)

すべてのコールバックはオーバーライド可能です:

  • allow_incoming_connection, allow_outgoing_connection: 接続バリデーション用の bool を返します
  • after_incoming_connection, after_outgoing_connection: 接続確立後のロジックを処理します
  • after_incoming_connection_removed, after_outgoing_connection_removed: 切断時の処理を行います
  • before_value_set: 値が設定される前に変更された値を返します
  • after_value_set: パラメータ値の変更に反応して処理を行います
  • validate_before_workflow_run, validate_before_node_run: list[Exception] | None を返します
  • on_griptape_event: ワークフローイベントを処理します
  • initialize_spotlight: スポットライト機能のセットアップを行います
  • get_next_control_output: 制御フロー用の Parameter | None を返します

ヘルパーメソッド (Helper Methods)

  • hide_parameter_by_name(), show_parameter_by_name()
  • append_value_to_parameter()
  • publish_update_to_parameter()
  • show_message_by_name(), hide_message_by_name(), get_message_by_name_or_element_id()

非同期 API 連携 (Asynchronous API Integration)

ノードがエンジンの動作を停止(ストール)させることなく長時間実行される処理を実行するには、2 つの方法があります:

  1. async def aprocess() をオーバーライドする(推奨)。 エンジンはそのイベントループ上でこれを await するため、真に非同期な連携(非同期 HTTP クライアント、await asyncio.sleep() によるポーリングなど)がエンジンの他の部分と並行して実行されます。
  2. process() をオーバーライドし、呼び出し可能オブジェクト(callable)を yield する (AsyncResult)。 yield された各呼び出し可能オブジェクトは バックグラウンドスレッド上で同期的に 実行されます — エンジンの応答性は維持されますが、処理自体はブロッキングかつ順次実行されるコードのままとなります。

新しい連携を作成する場合は、aprocess() と非同期 I/O を選択してください。書き換えたくない同期ライブラリ(requests、ブロッキング SDK など)をベースに連携が構築されている場合にのみ、yield パターンを使用してください。

aprocess() による非同期処理(推奨)

import asyncio

import httpx
from griptape_nodes.exe_types.node_types import ControlNode

POLLING_INTERVAL = 10  # 秒単位(API 推奨値を使用)
MAX_POLLING_ATTEMPTS = 60  # 最大 10 分間


class MyAsyncNode(ControlNode):
    async def aprocess(self) -> None:
        """リクエストを非同期に処理します。"""
        try:
            # 安全なデフォルト値を設定
            self._set_safe_defaults()

            # API キーを検証
            api_key = self._validate_api_key()

            async with httpx.AsyncClient(timeout=60) as client:
                # タスクを送信
                task_id = await self._submit_task(client, api_key)

                # 完了をポーリング
                result = await self._poll_for_completion(client, task_id, api_key)

            # 結果を処理
            self.parameter_output_values["output"] = result

        except Exception as e:
            self._set_safe_defaults()
            self._log(f"処理に失敗しました: {e}")
            raise RuntimeError(f"{self.name}: {e}") from e

    async def _submit_task(self, client: httpx.AsyncClient, api_key: str) -> str:
        response = await client.post(
            "https://api.example.com/v1/tasks",
            json=self._build_payload(),
            headers={"Authorization": f"Bearer {api_key}"},
        )
        response.raise_for_status()
        return response.json()["task_id"]

    async def _poll_for_completion(self, client: httpx.AsyncClient, task_id: str, api_key: str) -> str:
        for attempt in range(MAX_POLLING_ATTEMPTS):
            await asyncio.sleep(POLLING_INTERVAL)  # aprocess 内では絶対に time.sleep() を使用しないこと

            response = await client.get(
                "https://api.example.com/v1/query/task",
                params={"task_id": task_id},
                headers={"Authorization": f"Bearer {api_key}"},
            )
            response.raise_for_status()
            status_data = response.json()

            if status_data["status"] == "Success":
                return status_data["result"]
            if status_data["status"] == "Fail":
                error_msg = status_data.get("error_message", "不明なエラー")
                raise RuntimeError(f"タスクが失敗しました: {error_msg}")
            # "Processing", "Pending" などの場合はポーリングを継続

        raise RuntimeError(f"タスクが {MAX_POLLING_ATTEMPTS * POLLING_INTERVAL} 秒以内に完了しませんでした")

重要なポイント:

  • process() の代わりに async def aprocess() をオーバーライドします — エンジンがこれを直接 await します
  • 全体を通して非同期 I/O を使用します: リクエストには httpx.AsyncClient、ポーリング待機には await asyncio.sleep() を使用します
  • aprocess() 内でブロッキング呼び出し(requests、time.sleep())を行うと、エンジンのイベントループが停止します — やむを得ずブロッキング関数を呼び出す必要がある場合は、await asyncio.to_thread(blocking_fn) でラップしてください
  • ベースクラスのデフォルトの aprocess() は process() をラップしているため、ノードではどちらか一方のみをオーバーライドすれば十分です

バックグラウンドスレッドでのブロッキング処理 (process() + yield)

同期ライブラリ上に構築された連携の場合は、process() をオーバーライドして呼び出し可能オブジェクト(callable)を yield します。エンジンは yield された各呼び出し可能オブジェクトをバックグラウンドスレッド上で同期的に実行し、その戻り値でジェネレータを再開します — エンジンの応答性は保たれますが、これによって処理自体が非同期化される わけではありません:

from griptape_nodes.exe_types.node_types import ControlNode, AsyncResult


class MyBlockingNode(ControlNode):
    def process(self) -> AsyncResult | None:
        """ブロッキング処理をバックグラウンドスレッドに yield します。"""
        yield lambda: self._process()

    def _process(self) -> None:
        """メイン処理メソッド(バックグラウンドスレッド上で同期的に実行されます)。"""
        try:
            # 安全なデフォルト値を設定
            self._set_safe_defaults()

            # API キーを検証
            api_key = self._validate_api_key()

            # タスクを送信
            task_id = self._submit_task(api_key)

            # 完了をポーリング
            result = self._poll_for_completion(task_id, api_key)

            # 結果を処理
            self.parameter_output_values["output"] = result

        except Exception as e:
            self._set_safe_defaults()
            self._log(f"処理に失敗しました: {e}")
            raise RuntimeError(f"{self.name}: {str(e)}") from e

重要なポイント:

  • process() は AsyncResult | None を返し、呼び出し可能オブジェクトを yield します
  • yield された各呼び出し可能オブジェクトはバックグラウンドスレッド上で同期的に実行され、ジェネレータはその戻り値で再開されます
  • 既存の同期コード(requests、ブロッキング SDK)には適していますが、新しい連携を作成する場合は aprocess() を推奨します

長時間実行タスク向けのポーリングパターン

非同期タスク処理を使用する API(動画生成、モデルのトレーニングなど)と連携する場合は、3 ステップのパターンを実装します。以下の例では、前述のバックグラウンドスレッドパターンに適した同期的な requests 呼び出しを使用しています。aprocess() 内では、前述の通り代わりに httpx.AsyncClient と await asyncio.sleep() を使用してください。

ステップ 1: タスクの送信

def _submit_task(self, params: dict[str, Any], headers: dict[str, str]) -> dict[str, Any]:
    """タスクを送信し、task_id を含むレスポンスを返します。"""
    payload = self._build_payload(params)

    response = requests.post(self.API_BASE_URL, json=payload, headers=headers, timeout=DEFAULT_TIMEOUT)
    response.raise_for_status()

    response_data = response.json()
    task_id = response_data.get("task_id")
    return response_data

ステップ 2: ステータスのポーリング

POLLING_INTERVAL = 10  # 秒単位(API 推奨値を使用)
MAX_POLLING_ATTEMPTS = 60  # 最大 10 分間


def _poll_for_completion(self, task_id: str, headers: dict[str, str]) -> str | None:
    """タスク完了を API にポーリングし、結果識別子を返します。"""
    query_url = "https://api.example.com/v1/query/task"

    for attempt in range(MAX_POLLING_ATTEMPTS):
        time.sleep(POLLING_INTERVAL)  # 各ポーリング前に待機

        response = requests.get(
            query_url,
            headers=headers,
            params={"task_id": task_id},  # パスではなくクエリパラメータを使用
            timeout=DEFAULT_TIMEOUT,
        )
        response.raise_for_status()

        status_data = response.json()
        status = status_data.get("status")

        self._log(f"ポーリング試行 {attempt + 1}: ステータス = {status}")

        if status == "Success":
            file_id = status_data.get("file_id")
            return file_id
        elif status == "Fail":
            error_msg = status_data.get("error_message", "不明なエラー")
            raise RuntimeError(f"タスクが失敗しました: {error_msg}")
        # "Processing", "Pending" などの場合はポーリングを継続

    raise RuntimeError(f"タスクが {MAX_POLLING_ATTEMPTS * POLLING_INTERVAL} 秒以内に完了しませんでした")

ステップ 3: 結果の取得

def _retrieve_result(self, file_id: str, headers: dict[str, str]) -> str:
    """結果識別子からダウンロード URL を取得します。"""
    retrieve_url = "https://api.example.com/v1/files/retrieve"

    response = requests.get(retrieve_url, headers=headers, params={"file_id": file_id}, timeout=DEFAULT_TIMEOUT)
    response.raise_for_status()

    response_data = response.json()
    download_url = response_data.get("file", {}).get("download_url")

    return download_url

考慮すべき重要事項:

  • 常に API 推奨のポーリング間隔を使用してください(通常は 5〜10 秒)
  • 無限ループを防ぐため、妥当な最大試行回数を設定してください
  • task_id にはパスパラメータではなくクエリパラメータを使用してください(API ドキュメントで確認してください)
  • すべてのステータス状態(Success、Fail、Processing、Pending)を適切に処理してください
  • デバッグのためにポーリングの試行状況をログに出力してください
  • 失敗時には安全なデフォルト値を設定してください

入力に応じた動的エンドポイント選択

接続された入力に応じてノードが複数のモードで動作する場合(例: 画像が提供されている場合は image-to-video、提供されていない場合は text-to-video)、単一の URL をハードコードするのではなく、process メソッド内で動的に API エンドポイントを選択します:

IMAGE2VIDEO_URL = "https://api.example.com/v1/videos/image2video"
TEXT2VIDEO_URL = "https://api.example.com/v1/videos/text2video"


def _process(self):
    image_data = self._get_image_data("start_frame")
    has_images = image_data is not None

    if has_images:
        api_url = IMAGE2VIDEO_URL
    else:
        api_url = TEXT2VIDEO_URL

    payload = self._build_payload()
    if image_data:
        payload["image"] = image_data

    response = requests.post(api_url, headers=headers, json=payload, timeout=30)
    # ... ポーリングでもステータス確認に同じ api_url を使用します
    poll_url = f"{api_url}/{task_id}"

これにより、ユーザーがテキストのみの生成を行いたい場合に画像入力を要求することを回避し、各モードに対して正しい API エンドポイントが呼び出されるようになります。ポーリング URL にも同じベースエンドポイントを使用する必要があります。

画像アーティファクトの Base64 変換

以下の ImageArtifact 分岐は、過去に保存された古いワークフローや、まだ更新されていない上流ノードから渡される可能性のある値をサポートし続けるために用意されています。新しいパラメータでは、ImageArtifact ではなく ImageUrlArtifact(ParameterImage 経由)を宣言してください — 詳細は パラメータのペイロードサイズ を参照してください。

重要: localhost URL の処理

外部 API に画像を送信する際、静的ストレージからの ImageUrlArtifact の URL は localhost となり、外部サービスからはアクセスできません。必ず localhost の URL を検出して Base64 に変換してください:

import base64


def _get_image_data(self, image_artifact: ImageArtifact | ImageUrlArtifact) -> str:
    """画像アーティファクトを URL または Base64 データ URI に変換します。"""

    # ImageUrlArtifact - localhost かパブリック URL かをチェック
    if isinstance(image_artifact, ImageUrlArtifact):
        url = image_artifact.value

        # 外部 API の場合、localhost URL は Base64 に変換する必要がある
        if url.startswith(("http://localhost", "http://127.0.0.1", "https://localhost", "https://127.0.0.1")):
            self._log(f"localhost URL を Base64 に変換中: {url[:100]}...")
            response = requests.get(url, timeout=30)
            response.raise_for_status()
            image_bytes = response.content

            # ヘッダーから MIME タイプを検出
            mime_type = response.headers.get("content-type", "image/jpeg")
            if not mime_type.startswith("image/"):
                mime_type = "image/jpeg"

            base64_data = base64.b64encode(image_bytes).decode("utf-8")
            return f"data:{mime_type};base64,{base64_data}"

        # パブリック URL はそのまま渡すことができる
        self._log(f"パブリック URL を使用: {url[:100]}...")
        return url

    # ImageArtifact - .base64 プロパティを使用(推奨される方法)
    if isinstance(image_artifact, ImageArtifact):
        # 推奨: 組み込みプロパティを使用
        if hasattr(image_artifact, "base64") and hasattr(image_artifact, "mime_type"):
            base64_data = image_artifact.base64  # 生の Base64(プレフィックスなし)
            mime_type = image_artifact.mime_type  # 例: 'image/jpeg'

            # すでにデータ URI プレフィックスが付いているか確認
            if base64_data.startswith("data:"):
                self._log("ImageArtifact.base64 を使用(すでにデータ URI 形式)")
                return base64_data

            # データ URI プレフィックスを付与
            self._log(f"MIME タイプ {mime_type} で ImageArtifact.base64 を使用")
            return f"data:{mime_type};base64,{base64_data}"

        # フォールバック: 手動バイト抽出
        self._log("手動 Base64 エンコードにフォールバック中")
        if hasattr(image_artifact, "value") and hasattr(image_artifact.value, "read"):
            image_artifact.value.seek(0)
            image_bytes = image_artifact.value.read()
        elif hasattr(image_artifact, "data"):
            if isinstance(image_artifact.data, bytes):
                image_bytes = image_artifact.data
            elif hasattr(image_artifact.data, "read"):
                image_artifact.data.seek(0)
                image_bytes = image_artifact.data.read()
            else:
                raise ValueError("サポートされていない ImageArtifact 形式です")
        else:
            raise ValueError("サポートされていない ImageArtifact 形式です")

        # PIL で MIME タイプを検出
        mime_type = "image/jpeg"
        try:
            from PIL import Image
            from io import BytesIO

            img = Image.open(BytesIO(image_bytes))
            format_to_mime = {"JPEG": "image/jpeg", "PNG": "image/png", "WEBP": "image/webp"}
            mime_type = format_to_mime.get(img.format, "image/jpeg")
        except Exception:
            pass

        base64_data = base64.b64encode(image_bytes).decode("utf-8")
        return f"data:{mime_type};base64,{base64_data}"

    raise ValueError("サポートされていないアーティファクト型です")

重要なポイント:

  1. 常に localhost URL を検出する - 外部 API はこれらにアクセスできません
  2. ImageArtifact.base64 プロパティを使用する - Griptape の適切な方法です(生の Base64 を返します)
  3. ImageArtifact.mime_type プロパティを使用する - 自動的に MIME タイプを検出します
  4. どのパスが実行されたかをログに出力する - デバッグに不可欠です
  5. localhost のファイルをダウンロードする - API に送信する前に Base64 に変換します

パラメータ定義:

Parameter(
    name="image_input",
    input_types=["ImageUrlArtifact", "ImageArtifact"],  # ImageArtifact: レガシー入力互換性の目的のみ
    type="ImageUrlArtifact",
    tooltip="画像入力(ファイルまたは URL)",
    ui_options={"clickable_file_browser": True},  # ファイルブラウザを有効化
)

複数画像入力のバリデーション

ノードが複数の画像パラメータを受け付ける場合は、パラメータ名が明確に分かる再利用可能な検証メソッドを使用します。上記と同様に、ImageArtifact 分岐はレガシー入力の処理用であり、新規パラメータを ImageArtifact に対して宣言する理由にはなりません:

def _validate_image(self, image_artifact: ImageArtifact | ImageUrlArtifact, param_name: str) -> list[Exception]:
    """エラーメッセージにパラメータ名を含めて画像を検証します。"""
    exceptions = []

    if isinstance(image_artifact, ImageArtifact):
        # 画像バイト列を取得
        if hasattr(image_artifact, "value") and hasattr(image_artifact.value, "read"):
            image_artifact.value.seek(0)
            image_bytes = image_artifact.value.read()
            image_artifact.value.seek(0)
        else:
            return exceptions

        # サイズの検証
        size_mb = len(image_bytes) / (1024 * 1024)
        if size_mb >= 20:
            exceptions.append(ValueError(f"{self.name}: {param_name} のサイズは 20MB 未満である必要があります(現在: {size_mb:.1f}MB)"))

        # 形式と寸法の検証
        try:
            from PIL import Image
            from io import BytesIO

            img = Image.open(BytesIO(image_bytes))

            if img.format not in ["JPEG", "PNG", "WEBP"]:
                exceptions.append(
                    ValueError(f"{self.name}: {param_name} の形式は JPG、PNG、または WebP である必要があります(現在: {img.format})")
                )

            width, height = img.size
            short_edge = min(width, height)
            if short_edge <= 300:
                exceptions.append(
                    ValueError(f"{self.name}: {param_name} の短辺は 300px より大きい必要があります(現在: {short_edge}px)")
                )
        except ImportError:
            self._log("検証用の PIL が利用できません")
        except Exception as e:
            self._log(f"{param_name} の検証中にエラーが発生しました: {e}")

    return exceptions


def validate_before_node_run(self) -> list[Exception] | None:
    """すべての画像パラメータを検証します。"""
    exceptions = []

    # 各画像パラメータを個別に検証
    first_frame = self.get_parameter_value("first_frame_image")
    if first_frame:
        exceptions.extend(self._validate_image(first_frame, "first_frame_image"))

    last_frame = self.get_parameter_value("last_frame_image")
    if last_frame:
        exceptions.extend(self._validate_image(last_frame, "last_frame_image"))

    return exceptions if exceptions else None

メリット:

  • どの画像パラメータに問題があるかを特定できる明確なエラーメッセージ
  • 複数の画像入力にわたって再利用可能な検証ロジック
  • 各パラメータに対する個別の検証
  • ユーザーにとって分かりやすく対処可能なフィードバック

モデル依存のパラメータ管理

異なるモデルが異なるパラメータの組み合わせをサポートしている場合:

def after_value_set(self, parameter: Parameter, value: Any) -> None:
    """モデル依存のパラメータ表示および選択肢を処理します。"""
    if parameter.name == "model":
        if value == "AdvancedModel":
            # モデル固有のパラメータを表示
            self.show_parameter_by_name("advanced_option")

            # ドロップダウンの選択肢を動的に更新
            resolution_param = self.get_parameter_by_name("resolution")
            if resolution_param:
                for child in resolution_param.children:
                    if hasattr(child, "choices"):
                        child.choices = ADVANCED_MODEL_RESOLUTIONS
                        break
        else:
            # 他のモデルの場合は非表示にしてリセット
            self.hide_parameter_by_name("advanced_option")

            # 標準の選択肢に更新
            resolution_param = self.get_parameter_by_name("resolution")
            if resolution_param:
                for child in resolution_param.children:
                    if hasattr(child, "choices"):
                        child.choices = STANDARD_RESOLUTIONS
                        break
                self.set_parameter_value("resolution", "720P")

    return super().after_value_set(parameter, value)

モデル固有のバリデーション:

def validate_before_node_run(self) -> list[Exception] | None:
    """モデル固有のパラメータの組み合わせを検証します。"""
    exceptions = []

    model = self.get_parameter_value("model")
    duration = self.get_parameter_value("duration")
    resolution = self.get_parameter_value("resolution")

    # 例: 10秒の再生時間は特定のモデル/解像度でのみサポートされる
    if duration == 10:
        if model != "AdvancedModel":
            exceptions.append(ValueError(f"{self.name}: 10秒の再生時間は AdvancedModel でのみサポートされています"))
        elif resolution == "4K":
            exceptions.append(ValueError(f"{self.name}: 10秒の再生時間は 4K 解像度ではサポートされていません"))

    # モデル固有の必須パラメータ要件
    if model in ["ModelB", "ModelC"]:
        required_param = self.get_parameter_value("required_for_model_b_c")
        if not required_param:
            exceptions.append(ValueError(f"{self.name}: {model} にはパラメータが必要です"))

    return exceptions if exceptions else None

非推奨モデルの移行とユーザー通知

モデルプロバイダーがエンドポイントを非推奨にする場合(例: プレビューモデルが GA 正式版に置き換わる場合)、ノードは保存されたワークフローを自動的に移行しながらユーザーに通知する必要があります。このパターンでは、連携して機能する 3 つのコンポーネントを使用します:

  1. 古いモデル名と置換先をマッピングする DEPRECATED_MODELS 辞書
  2. 閉じることができる情報バナーとして機能する非表示の ParameterMessage 要素
  3. 非推奨の値が適用される前にインターセプトして置換するための before_value_set ライフサイクルフック

ステップ 1: 非推奨マップと現在のモデルを定義する

from griptape_nodes.exe_types.core_types import Parameter, ParameterMessage
from griptape_nodes.traits.button import Button

MODELS = [
    "veo-3.1-generate-001",
    "veo-3.1-fast-generate-001",
]

# 非推奨のモデル名と置換先のモデル名のマッピング。
# 保存されたワークフローがこれらを参照している場合、ノードは自動移行します。
DEPRECATED_MODELS: dict[str, str] = {
    "veo-3.1-generate-preview": "veo-3.1-generate-001",
    "veo-3.1-fast-generate-preview": "veo-3.1-fast-generate-001",
    "veo-3.0-generate-001": "veo-3.1-generate-001",
    "veo-2.0-generate-001": "veo-3.1-generate-001",
}

ステップ 2: __init__ で非表示の ParameterMessage を追加する

UI 上でモデルセレクターの近くに表示されるよう、model パラメータの後に追加します。hide=True により、必要になるまで非表示のままになります。on_click を持つ Button トレイトによって、ユーザーに「Dismiss(閉じる)」ボタンを提供します。

def __init__(self, **kwargs):
    super().__init__(**kwargs)

    # ... 上記で model パラメータが追加されている前提 ...

    # 非推奨通知メッセージ(非表示) — 非推奨モデルが検出されたときに表示
    self.add_node_element(
        ParameterMessage(
            name="model_deprecation_notice",
            title="モデル非推奨通知",
            variant="info",
            value="",
            traits={
                Button(
                    full_width=True,
                    on_click=lambda _, __: self.hide_message_by_name("model_deprecation_notice"),
                )
            },
            button_text="閉じる",
            hide=True,
        )
    )

ステップ 3: before_value_set を実装して非推奨モデルをインターセプトする

before_value_set はパラメータの値が適用される前に呼び出されます。after_value_set(およびモデル値に依存するロジック)は置換後の値を受け取ることになるため、非推奨モデルを後継モデルに置き換えるにはここが最適な場所です。

def before_value_set(self, parameter: Parameter, value: Any) -> Any:
    """非推奨モデルを自動移行し、非推奨通知を表示します。"""
    if parameter.name == "model" and value in DEPRECATED_MODELS:
        replacement = DEPRECATED_MODELS[value]
        message = self.get_message_by_name_or_element_id("model_deprecation_notice")
        if message is not None:
            message.value = (
                f"モデル '{value}' は非推奨となりました。"
                f"モデルは自動的に '{replacement}' に更新されました。"
                "この変更を適用するためにワークフローを保存してください。"
            )
            self.show_message_by_name("model_deprecation_notice")
        value = replacement

    return super().before_value_set(parameter, value)

ステップ 4: ユーザーが有効なモデルを選択したときに通知を非表示にする

after_value_set で、現在のモデルが非推奨でない場合にバナーを閉じます。これにより、ユーザーが移行後に手動で別の有効なモデルを選択した場合に対応できます。

def after_value_set(self, parameter: Parameter, value: Any) -> None:
    if parameter.name == "model":
        # ... モデル固有のロジック(duration の選択肢の更新など) ...
        if value not in DEPRECATED_MODELS:
            self.hide_message_by_name("model_deprecation_notice")

    return super().after_value_set(parameter, value)

エンドツーエンドの動作の流れ:

  1. ユーザーが "veo-3.1-generate-preview" で保存されたワークフローを開きます。
  2. フレームワークは保存された値を使って before_value_set を呼び出します。
  3. フックはそれが DEPRECATED_MODELS に存在することを検出し、"veo-3.1-generate-001" に置き換えて、情報バナーを表示します。
  4. 置換された値で after_value_set が発火します — 有効な GA モデルが渡されるため、モデル依存の UI 更新(再生時間の選択肢、パラメータの可視性など)が正常に機能します。
  5. ユーザーには次のバナーが表示されます: 「モデル 'veo-3.1-generate-preview' は非推奨となりました。モデルは自動的に 'veo-3.1-generate-001' に更新されました。この変更を適用するためにワークフローを保存してください。」
  6. ユーザーはバナーを閉じることができます。また、次に有効なモデルが選択されたときに自動的に非表示になります。

使用される主要な API メソッド:

メソッド 役割
self.add_node_element(ParameterMessage(...)) メッセージ要素をノードに追加します
self.get_message_by_name_or_element_id(name) 値を更新するためにメッセージ要素を取得します
self.show_message_by_name(name) 非表示のメッセージを表示状態にします
self.hide_message_by_name(name) メッセージを再び非表示にします

リファレンス実装:

  • griptape_nodes_library/config/prompt/griptape_cloud_prompt.py 内の GriptapeCloudPrompt(標準ライブラリ)
  • griptape-nodes-library-googleai 外部ライブラリ内の VeoVideoGenerator、VeoImageToVideoGenerator、VeoTextToVideoWithRef

API 連携のための高度なデバッグログ

外部 API と連携するノードでは、問題を迅速に診断できるように包括的なデバッグログを実装します:

# タスク送信 - 完全なレスポンスをログ出力
def _submit_task(self, params: dict, headers: dict) -> dict:
    response = requests.post(API_URL, json=payload, headers=headers)
    response.raise_for_status()

    response_data = response.json()
    self._log(f"タスク送信レスポンス: {json.dumps(response_data, indent=2)}")
    return response_data


# ペイロードサイズ - 送信前にデータサイズをログ出力
def _log_request(self, payload: dict) -> None:
    if "first_frame_image" in payload:
        img_len = len(payload.get("first_frame_image", ""))
        self._log(f"first_frame_image データ長: {img_len} 文字 (~{img_len / 1024:.1f}KB)")

    if "last_frame_image" in payload:
        img_len = len(payload.get("last_frame_image", ""))
        self._log(f"last_frame_image データ長: {img_len} 文字 (~{img_len / 1024:.1f}KB)")


# エラーレスポンス - API エラーの詳細を完全にログ出力
def _poll_for_completion(self, task_id: str, headers: dict) -> str:
    status_data = response.json()
    status = status_data.get("status")

    if status == "Fail":
        # デバッグのために完全なエラーレスポンスをログ出力
        self._log(f"完全な API エラーレスポンス: {json.dumps(status_data, indent=2)}")
        error_msg = status_data.get("error_message", "不明なエラー")
        raise RuntimeError(f"タスクが失敗しました: {error_msg}")


# 処理パス - どのコードパスが実行されたかをログ出力
def _get_image_data(self, image_artifact) -> str:
    if isinstance(image_artifact, ImageUrlArtifact):
        if url.startswith("http://localhost"):
            self._log(f"localhost URL を Base64 に変換中: {url[:100]}...")
        else:
            self._log(f"パブリック URL を使用: {url[:100]}...")
    elif isinstance(image_artifact, ImageArtifact):
        if hasattr(image_artifact, "base64"):
            self._log(f"MIME タイプ {mime_type} で ImageArtifact.base64 を使用")
        else:
            self._log("手動 Base64 エンコードにフォールバック中")

ログに記録すべき項目:

  • 完全な API レスポンス(送信、ポーリング、取得)
  • ペイロードサイズ(特に Base64 データ)
  • 処理パス(どのコード分岐が実行されたか)
  • 使用されている モデル/パラメータの組み合わせ
  • エラーの詳細(API からの完全なエラーレスポンス)

メリット:

  • どこで障害が発生したかを素早く特定可能
  • どのようなデータが送信されているかを把握可能
  • どのコードパスが実行されたかを追跡可能
  • 正確な API エラーメッセージとコードを取得可能
  • 問題を再再現することなくデバッグ可能

API ドキュメントの検証

重要なベストプラクティス: API 仕様は常に公式ドキュメントを直接確認して検証してください。

回避すべき一般的な落とし穴:

  1. モデル名: 大文字小文字の完全一致を確認(video-01 ではなく MiniMax-Hailuo-02)
  2. エンドポイント: 正確な URL を確認(/v1/video_generation/{id} ではなく /v1/query/video_generation)
  3. パラメータ: クエリパラメータとパスパラメータの違いを確認
  4. レスポンス構造: 正確なフィールド名を確認(file_list ではなく file_id)
  5. ポーリング間隔: API 推奨値を使用

例: 正しいポーリングと誤ったポーリング:

# ✅ 正しい: クエリパラメータ
response = requests.get("https://api.example.com/v1/query/task", params={"task_id": task_id})

# ❌ 誤り: パスパラメータ(API が明示的にこれを指定していない限り)
response = requests.get(f"https://api.example.com/v1/query/task/{task_id}")

ドキュメントにアクセスできない場合:

  • Web ページにアクセスできないこと(JavaScript を多用したドキュメントなど)を明確に説明する
  • ユーザーに関連するドキュメント箇所の提供を求める
  • 検証なしに API パターンを推測・想定しない
  • サンプルコードが提供された場合は、それに基づいて実装を更新する