PipelineDocs
GitHub
Docs / プラグイン開発

プラグインSDK

プラグインは python_module executor から呼ばれる Python モジュールです。 3関数 (setup / process / cleanup) を export するのが基本契約で、 あとは plugin.yaml でフォーム項目とメタ情報を宣言するだけです。

ディレクトリ構成

plugins/
  my_plugin/
    plugin.yaml         # マニフェスト (必須)
    main.py             # setup / process / cleanup を export (module名は任意)
    requirements.txt    # (任意) pip 依存
    web/
      panel.html        # (任意) ui_panel: true のときのカスタムUI
    README.md           # (任意)

プラグインルートは環境変数 PIPELINE_PLUGIN_ROOT (既定 /opt/pipeline/plugins) です。ここに置いたディレクトリが GET /api/v1/plugins/available で列挙され、workload 作成フォームのプラグインドロップダウンに出ます。

plugin.yaml マニフェスト

name: my-plugin
description: |
  ここにプラグインの説明を書きます。Web UI に表示されます。

# workload 作成フォームに出る入力項目 (setup() の kwargs になる)
init_kwargs:
  - key: api_url
    type: str
    default: https://api.example.com
    label: API endpoint
    help: タスクが叩く API の URL
    required: true

  - key: timeout_secs
    type: float
    default: 10
    label: HTTP timeout
    min: 1
    max: 300

  - key: mode
    type: enum
    default: fast
    label: 動作モード
    options: [fast, safe]

  - key: dry_run
    type: bool
    default: false
    label: 書き込みをスキップ

  - key: api_key
    type: secret
    label: API key

# フォームには出さず、raw JSON で編集する項目 (慣習: 制御プレーン URL 等)
hidden_kwargs:
  - workload_slug
  - control_url

# (任意) カスタム UI パネル
ui_panel: true
ui_panel_mode: image        # video (既定) / image

# (任意) ループ型プラグインの宣言
workload_type: self_loop    # 省略時は通常の per-task 型

マニフェストのトップレベルキー

キー意味
namestrプラグイン表示名。
descriptionstrWeb UI に出す説明。
init_kwargslistフォーム項目。各値は setup() の kwargs になります。
hidden_kwargslist[str]フォームから隠す項目名 (raw JSON で直接編集する想定)。
ui_panelbooltrue で専用 UI パネルを有効化。
ui_panel_modestrパネルに渡す ?mode=video (既定) / image
workload_typestr/nullself_loop でループ型宣言 (watchdog 対象)。省略/event_driven は通常型。

init_kwargs エントリのフィールド

フィールド意味
keystr (必須)setup() に渡す kwarg 名。
typestr (必須)int / float / str / path / bool / enum / secret
defaultanyフォームの初期値。
label / helpstrフォームのラベル / 補助説明。
min / maxnumber数値の下限 / 上限。
optionslistenum の選択肢。
requiredbool必須フラグ。

type と Web UI の input 対応

typeWeb UI の input
str / pathTextInput (path は path 用 placeholder)
int / floatNumberInput (整数 / 小数)
boolCheckbox
enumSelect (options 必須)
secretPasswordInput
hidden_kwargs は自動注入ではありません: hidden_kwargs はあくまで「フォームに出さない」UI ヒントです。制御プレーンが workload_slug / control_url を自動で入れてくれる保証はないので、プラグイン側は必ずフォールバック付きで読みます (例: kwargs.get("control_url") or "http://localhost:8000")。実際に使う場合は workload の init_kwargs に raw JSON で値を入れてください。

main.py — プラグイン契約

from __future__ import annotations
from typing import Any


def setup(**kwargs: Any) -> Any:
    """プロセス起動時に1回だけ呼ばれます (最初の claim 時)。
    重い初期化 (モデルロード / DB接続) はここで。戻り値が state になり、
    process() / cleanup() に渡されます。省略可 (その場合 state=None)。
    ここで例外を投げると executor build が失敗し、claim済タスクは失敗扱いに。"""
    return {
        "api_url": kwargs["api_url"],
        "timeout": float(kwargs.get("timeout_secs", 10)),
        "client": _make_http_client(),
    }


def process(task: Any, ctx: Any, state: Any) -> dict[str, Any] | None:
    """1タスク毎に呼ばれます (必須)。
    - task  … pk / extra / attempt / workload_slug を持つ Task
    - ctx   … deadline / workdir / env / workload_config を持つ ExecutionContext
    - state … setup() の戻り値 (キーワード引数 state= で渡される)
    戻り dict は runs.output_json に保存されます。例外を上げると失敗扱いです。"""
    r = state["client"].get(f"{state['api_url']}/{task.pk}", timeout=state["timeout"])
    return {"status_code": r.status_code, "size": len(r.content)}


def cleanup(state: Any) -> None:
    """プロセス停止 / 設定変更 / filter変更 / eviction時に呼ばれます (任意)。
    接続クローズなどを行います。ここでの例外は握り潰されます。"""
    if state:
        state["client"].close()
重要 — 引数の渡し方: worker は process(task, ctx, state=state) のように stateキーワード引数で渡します。第3引数の名前は必ず state にしてください。setup()setup(**init_kwargs) で呼ばれるので、init_kwargs の各キーがそのままキーワード引数になります。

task / ctx / 戻り値

task (dataclass)process() の第1引数。

属性意味
pkAnyタスクの主キー。
workload_slugstr所属 workload の slug。
attemptint何回目の試行か (1始まり)。
extradictenqueue 時に添えた JSON。

ctx (ExecutionContext) — 第2引数。

属性意味
deadlinedatetimelease 失効時刻。長時間処理はこれを目安に切り上げます。
workdirPathタスク専用の一時ディレクトリ。処理後に自動削除されます。
envdictos.environ のコピー。
workload_configdictexecutor_config のコピー。

process() の戻り値の扱い

戻り値結果
dict (キー _error なし)成功。output_json に保存。
dict{"_error": "..."} を含む失敗 (error = その文字列)。例外を投げずに失敗を返せます。
None成功 (output_json なし)。
True / False成功フラグとしてそのまま使用。
例外を raise失敗 (traceback / stderr が記録され、attempt 増加)。

モジュール解決 (module / source_path / callable)

python_moduleexecutor_config は3つの鍵でプラグインを特定します。

意味
source_pathプラグインのディレクトリ。dev は plugins/url_check、本番は /opt/pipeline/plugins/url_check 等。
moduleそのディレクトリ内の .py ファイル名 (拡張子なし)。__init__.py を持つサブパッケージ名でも可。
callable呼ぶ関数名。既定 process

内部では source_pathsys.path 先頭に加わり (from .lib import x のような相対 import が効く)、moduleプラグイン slug 込みのユニークな sys.modules キーでロードされます。

同名モジュールの衝突に注意 (と、その対策): 2つのプラグインが両方 dispatch_main.py のような同名ファイルを持つと、素の import では sys.modules キャッシュを共有し、片方のコードが両方の workload で動いてしまう事故が起きます (実際に発生しました)。Pipeline は source_path 指定時に spec_from_file_location + slug 入りキーでロードしてこれを物理的に隔離します。つまり ファイル名は他プラグインと被っても安全ですが、可読性のため main.py 等の分かりやすい名前を推奨します。

依存パッケージ (requirements.txt)

requirements.txt は「このプラグインが必要とする依存」を書くためのドキュメントです。Pipeline はこれを自動 install しません。プラグインが外部ライブラリを使うなら、ワーカーの venv に事前に入っている必要があります。方法は3つ。

  • ワーカーの venv に手で pip install -r requirements.txt
  • Deploysetup_commandpip install を書いて配信時に流す。
  • fleet 共通の venv を使い、依存をそこへ集約する。

標準ライブラリだけで書ければ依存ゼロで最も楽です (チュートリアルの url_checkurllib のみ)。

エラー処理とリトライ設計

  • 失敗の表明は2通り。 一時的な障害 (ネットワーク・5xx・タイムアウト) は例外を raise して max_attempts のリトライに乗せます。恒久的にダメな入力 (不正 URL・存在しない ID 等) は {"_error": "..."} を返して早々に確定させ、無駄なリトライを避けます。
  • 冪等性を前提に書く。 lease 失効やリトライで、同じ pk が2回以上 process() されることがあります。出力の書き込みは upsert / INSERT IGNORE にする、成果物のキーに pk を使う等で、二重実行しても壊れないようにします。
  • 時間を意識する。 長い処理は ctx.deadline を目安に切り上げます。max_duration_ms を設定すると超過を自動で失敗にできます。
  • OOM の自己申告。 エラー文言に CUDA OOM 系のシグネチャが含まれると、ワーカーが VRAM ピークを自己申告して以後の同時実行を絞ります。

ロギング

print(...)logging の出力はワーカーの stdout/stderr に出て、ServiceLog 画面に集約されます。プラグインでは logging.getLogger("pipeline.plugin.url-check") のように名前を付けると追いやすくなります。例外は自動で traceback が runs と ServiceLog に残ります。

開発ループとローカルテスト

process() の実体はワーカープロセス内にキャッシュされます (重い setup() を毎回走らせないため)。したがってコードを直したらワーカーを再起動してください。dev は pipeline run --dev を Ctrl-C して起動し直すだけ、本番は該当ホストの worker を restart し、__pycache__/*.pyc を消してから上げ直します。

fleet に載せる前に、process() は普通の関数として単体で叩けます。

python3 - <<'PY'
from types import SimpleNamespace
import main                       # plugins/url_check/ で実行
state = main.setup(timeout_secs=5)
task  = SimpleNamespace(pk="https://example.com", extra={}, attempt=1, workload_slug="url-check")
ctx   = SimpleNamespace(deadline=None, workdir="/tmp", env={}, workload_config={})
print(main.process(task, ctx, state))
PY

手順を通しで体験するには チュートリアルのステップ4〜5 をどうぞ。

バッチ処理 (任意)

process_batch(tasks, ctx, state) を export すると、worker は batch_size に応じて複数タスクをまとめて渡します。戻り値の list 長は len(tasks) と一致させてください (不一致は全件失敗)。GPU 推論など、まとめた方が効率が良い plugin 向けです。

def process_batch(tasks: list, ctx: Any, state: Any) -> list[dict]:
    embeddings = state["model"].encode([t.extra["text"] for t in tasks])
    return [{"dim": len(e)} for e in embeddings]   # len == len(tasks)

self_loop (ループ型) プラグイン

キューを消費する代わりに「回り続ける」プラグイン (dispatcher / 投入系 / supervisor) を作れます。workload_type: self_loop を宣言した上で、次の3点を実装します。

  1. ブートストラップsetup() の最後に、自分のキューへ最初の tick を HTTP 投入:
    POST {control_url}/api/v1/workloads/{slug}/tasks{"pk": "tick-1-<ts>"}
  2. 長い tickprocess() の1回が実作業。ループ型は ctx.deadline (= lease) の少し手前まで while で回し、ポーリング型は1パス実行して time.sleep(interval_s)
  3. 次 tick の自己 enqueueprocess() の最後に、次の tick pk を再び自己投入してループを継続します。

worker はこれを通常の per-task plugin と同じく process(task, ctx, state) として駆動します (task は合成 tick)。

watchdog と singleton: worker が tick 完了後・自己 enqueue 前に死ぬとループが止まります。これは制御プレーンの SelfLoopWatchdog が救済します (60秒毎に self_loop workload を走査し、pending=claimed=0 が300秒続いたらブートストラップ tick を1件投入)。また self_loop は多重起動すると増殖するため、max_concurrent_total: 1 + supervisor_enabled: false で単一化するのが定石です。

ストレージの使い方

MinIO / DB を使うプラグインは、接続の生成・認証・死活監視を Storage SDK (pipeline.storage) に任せられます。setup() で接続レジストリを作り、process() で名前指定のハンドルを使います。

サンプルと入手

最小サンプルが同梱されています。新規 plugin はこれをコピーするのが早いです。

Pipeline — GUI-first batch fleet · ぱっぷすラボ GitHub · ホーム