PocketBaseで軽量なワークフロー状態管理を実現する — Temporal/Airflowを使わない選択肢

PocketBaseで軽量なワークフロー状態管理を実現する — Temporal/Airflowを使わない選択肢

TemporalやAirflowほどの規模は不要だがExcel管理では限界、という場面でPocketBaseをワークフロー状態管理に採用しました。axiosインターセプターによる自動起動・認証やpb_hooksによるカスタムエンドポイント追加など、実装のポイントを紹介します。
2026.07.27

はじめに

機器交換業務を自動化する際、各案件の処理状態(申請受付、出荷待ち、交換完了など)を追跡する必要がありました。

TemporalやAirflowのような本格的なワークフローエンジンは導入コストが高く、Excelでの管理は限界がありました。そこでPocketBase(組み込み型バックエンド)をワークフローの状態管理に採用しました。

前提・環境

  • PocketBase(SQLiteベースの組み込みバックエンド)
  • Node.js + TypeScript + axios
  • PocketBaseカスタムフック(pb_hooks/

なぜPocketBaseか

選択肢 利点 この用途での課題
Excel/スプレッドシート 導入簡単 同時編集の競合、API連携困難
Temporal/Airflow 本格的なワークフロー管理 インフラ運用が重い、学習コスト
RDB + 自前API 柔軟 APIの実装コスト
PocketBase 組み込み型、REST API付き、カスタムフック 大規模には不向き

PocketBaseは単一バイナリで起動でき、SQLiteベースのデータベース + REST API + 管理UIがすぐ使えます。さらにpb_hooks/ディレクトリにJavaScriptファイルを置くだけで、カスタムエンドポイントやサーバーサイドロジックを追加できます。これが今回の選定で最も決め手になった機能です。

実装

全体フロー

以下の図は、axiosインターセプターによる自動起動と認証の流れを示しています。

pocketbase-lightweight-workflow-state-machine-interceptor-flow

PocketBaseサーバーの自動起動

アプリケーション起動時にPocketBaseが動いていなければ自動起動します。

Pocketbase.ts
class PocketBase {
  private api: AxiosInstance;

  constructor() {
    this.api = axios.create({
      baseURL: "http://127.0.0.1:8090",
    });

    // axiosインターセプターで自動接続チェック + 認証
    this.api.interceptors.request.use(async (config) => {
      await this.checkConnection();
      config.headers.Authorization = `Bearer ${await this.getAuthToken()}`;
      return config;
    });
  }

  private async checkConnection(): Promise<void> {
    const isRunning = await this.isServerRunning();
    if (isRunning) return;
    await this.initPocketBase();
  }

PocketBaseバイナリの起動は、プロセスのstdout/stderrを監視して"Server started at"が出力されたら完了とみなします。

Pocketbase.ts
  public initPocketBase(): Promise<void> {
    return new Promise<void>((resolve, reject) => {
      if (environment.CURRENT_OS === "darwin") {
        execSync(`chmod +x "${pocketbasePath}"`);
        execSync(`xattr -d com.apple.quarantine "${pocketbasePath}"`);
      }

      const process = spawn(pocketbasePath, ["serve"]);

      const timeout = setTimeout(() => {
        process.kill();
        reject(new Error("PocketBase startup timed out"));
      }, 10000);

      const onData = (chunk: Buffer) => {
        if (chunk.toString().includes("Server started at")) {
          clearTimeout(timeout);
          process.stdout.removeListener("data", onData);
          process.stderr.removeListener("data", onData);
          resolve();
        }
      };

      process.stdout.on("data", onData);
      process.stderr.on("data", onData);

      process.on("error", (err) => {
        clearTimeout(timeout);
        reject(err);
      });
    });
  }

macOSではxattr -d com.apple.quarantineでGatekeeperの隔離属性を解除する必要がありました。

axiosインターセプターでの認証管理

すべてのAPIリクエストに認証トークンを自動付与します。

Pocketbase.ts
  private async getAuthToken(): Promise<string> {
    // インターセプター内でのループを避けるため、axiosを直接使用
    const response = await axios.post(
      `${this.pocketbaseUrl}/api/admins/auth-with-password`,
      {
        identity: credentials.POCKETBASE_USERNAME,
        password: credentials.POCKETBASE_PASSWORD,
      }
    );
    return response.data.token;
  }

注意点として、getAuthToken()ではインスタンスのthis.apiではなくグローバルのaxiosを使っています。this.apiを使うとインターセプターが再帰的に呼ばれて無限ループになるためです。

ワークフロー操作

以下の図は、機器交換ワークフローの状態遷移を示しています。

pocketbase-lightweight-workflow-state-machine-state-transition

機器交換の状態管理に使うCRUD操作です。

Pocketbase.ts
  // ステータスでフィルタリング(カスタムエンドポイント)
  public async filterByStatus(
    status: string,
    equipmentType: string,
    daysAfter: number = 0
  ): Promise<EquipmentExchange[]> {
    const response = await this.api.post(
      "/api/collections/EquipmentExchange/filterByStatus",
      { status, equipmentType, daysAfter }
    );
    return response.data.items;
  }

  // ステータス更新
  public async updateStatus(serialNumber: string, status: string): Promise<void> {
    await this.api.post("/api/collections/EquipmentExchange/updateStatus", {
      serialNumber,
      status,
    });
  }

  // 重複チェック付きの一括登録
  public async registerEntries(entries: EquipmentExchange[]): Promise<void> {
    for (const entry of entries) {
      if (await this.checkForDuplicate(entry.serialNumber)) {
        Logger.error(`Duplicate found for ${entry.serialNumber}, skipping.`);
        continue;
      }
      await this.api.post("/api/collections/EquipmentExchange/records", entry);
    }
  }

カスタムフックによるサーバーサイドロジック(pb_hooks)

PocketBaseの最大の強みはpb_hooks/ディレクトリです。ここにJSファイルを置くだけで、PocketBaseが自動的に読み込み、カスタムエンドポイントとして公開します。

エンドポイントの定義パターン

routerAdd()でHTTPメソッド、パス、ハンドラを指定します。

pb_hooks/EquipmentExchange.pb.js
// ステータス + 機器種別 + 日付条件で複合フィルタリング
routerAdd("POST", "api/collections/EquipmentExchange/filterByStatus", (c) => {
  const status = $apis.requestInfo(c).data.status;
  const equipmentType = $apis.requestInfo(c).data.equipmentType;
  const daysAfter = $apis.requestInfo(c).data.daysAfter;
  const systemTime = $apis.requestInfo(c).data.systemTime;

  // DynamicModelで結果のスキーマを定義
  const records = arrayOf(
    new DynamicModel({
      serialNumber: "",
      accountId: "",
      contactName: "",
      shipmentAddress: "",
      deliveryDate: "",
      processStatus: "",
      equipmentType: "",
      trackingNumber: "",
      requestDate: "",
      lastCheckDate: "",
      retry: "",
    })
  );

  // 生SQLで複雑なフィルタリング
  const query = [
    `SELECT *`,
    `FROM EquipmentExchange`,
    `WHERE processStatus = '${status}'`,
    `AND equipmentType = '${equipmentType}'`,
    `AND retry < 3`,
    daysAfter === 0 ? "" : `AND '${systemTime}' >= deliveryDate`,
    daysAfter === 0
      ? ""
      : `AND '${systemTime}' >= datetime(requestDate, '+${daysAfter} days')`,
    `ORDER BY created ASC`,
  ].join(" ");

  $app.dao().db().newQuery(query).all(records);

  return c.json(200, { items: records });
});

ポイントは3つあります。

$apis.requestInfo(c).data でリクエストボディのJSONフィールドを取得します。PocketBase独自のAPIで、Express/Fastifyのreq.bodyに相当します。

DynamicModel でSQLの結果をマッピングするスキーマを定義します。PocketBaseのGoランタイム内でJSを実行するため、通常のORMは使えません。代わりにDynamicModelで列名と初期値を指定し、arrayOf()で配列として受け取ります。

$app.dao().db().newQuery(query).all(records) でSQLを実行します。PocketBaseの標準REST APIでは表現できない複雑なクエリ(複数条件のAND、日付計算、リトライ回数制限など)をサーバー側で処理できます。

更新系エンドポイント

ステータス更新やフィールド更新も同じパターンです。

pb_hooks/EquipmentExchange.pb.js
// ステータス更新
routerAdd("POST", "api/collections/EquipmentExchange/updateStatus", (c) => {
  const serialNumber = $apis.requestInfo(c).data.serialNumber;
  const status = $apis.requestInfo(c).data.status;
  const query = [
    `UPDATE EquipmentExchange`,
    `SET processStatus = '${status}'`,
    `WHERE serialNumber = '${serialNumber}'`,
  ].join(" ");

  $app.dao().db().newQuery(query).execute();

  return c.json(200, { message: "query executed successfully." });
});

// リトライカウンターのインクリメント
routerAdd("POST", "api/collections/EquipmentExchange/incRetry", (c) => {
  const serialNumber = $apis.requestInfo(c).data.serialNumber;
  const query = [
    `UPDATE EquipmentExchange`,
    `SET retry = retry + 1`,
    `WHERE serialNumber = '${serialNumber}'`,
  ].join(" ");

  $app.dao().db().newQuery(query).execute();

  return c.json(200, { message: "query executed successfully." });
});

// レコード削除
routerAdd("POST", "api/collections/EquipmentExchange/deleteRecord", (c) => {
  const serialNumber = $apis.requestInfo(c).data.serialNumber;
  const query = `DELETE FROM EquipmentExchange WHERE serialNumber = '${serialNumber}'`;

  $app.dao().db().newQuery(query).execute();

  return c.json(200, { message: "query executed successfully." });
});

SELECT系は.all(records)、UPDATE/DELETE系は.execute()を使い分けます。

複雑なバッチ更新

CTE(Common Table Expression)を使った集約更新もサーバー側で実行できます。

pb_hooks/EquipmentExchange.pb.js
// 同一配送先・同一日付のエントリをグループ化してバッチ処理用にマーク
routerAdd("POST", "api/collections/EquipmentExchange/groupEntries", (c) => {
  const records = arrayOf(
    new DynamicModel({
      shipmentAddress: "",
      deliveryDate: "",
      entryCount: "",
    })
  );

  const query = [
    `WITH GroupedEntries AS (`,
    `  SELECT shipmentAddress, deliveryDate, COUNT(*) as entryCount`,
    `  FROM EquipmentExchange`,
    `  WHERE processStatus = 'INIT'`,
    `  GROUP BY shipmentAddress, deliveryDate`,
    `  HAVING COUNT(*) > 1`,
    `)`,
    `UPDATE EquipmentExchange`,
    `SET processStatus = 'BATCH_READY'`,
    `FROM GroupedEntries`,
    `WHERE EquipmentExchange.shipmentAddress = GroupedEntries.shipmentAddress`,
    `  AND EquipmentExchange.deliveryDate = GroupedEntries.deliveryDate`,
    `  AND EquipmentExchange.processStatus = 'INIT';`,
  ].join(" ");

  $app.dao().db().newQuery(query).all(records);

  return c.json(200, { items: records });
});

このようなCTEを含む複雑なSQLも、PocketBaseのカスタムフック内で直接実行できます。SQLiteがサポートする構文であれば何でも使えるため、アプリケーション側で複数回のAPI呼び出しを行う代わりに、1回のリクエストで完結させられます。

操作ログの管理

機器交換の各ステップのログをPocketBaseに記録し、後から追跡できるようにしています。

Pocketbase.ts
  public async writeLog(log: ExchangeLog): Promise<void> {
    await this.api.post("/api/collections/ExchangeLogs/records", log);
  }

  public async deleteApiLogs(entries: { id: string }[]): Promise<void> {
    const deletePromises = entries.map((record) =>
      this.api.delete(`/api/collections/ExchangeLogs/records/${record.id}`)
    );
    await Promise.all(deletePromises); // 並列削除
  }

ログの削除はPromise.all()で並列化しています。

ハマったポイント

Bearer トークンのフォーマット: 初期実装でBearer${token}(スペースなし)としていたため、認証がサイレントに失敗していました。Bearer ${token}(スペースあり)が正しいフォーマットです。エラーメッセージが出ないため発見が遅れました。

macOSのGatekeeper: PocketBaseバイナリがダウンロードされたファイルとしてマークされ、実行がブロックされます。xattr -d com.apple.quarantineで隔離属性を解除する必要がありました。

まとめ

PocketBaseを使って、機器交換ワークフローの状態管理を実現しました。

  • 単一バイナリで起動、REST APIと管理UIが即利用可能
  • axiosインターセプターで自動起動 + 認証を透過的に処理
  • pb_hooks/のJSファイルでカスタムエンドポイントを追加(CTE含む複雑なSQLも実行可能)
  • DynamicModel + routerAdd()で、Go実装なしにサーバーサイドロジックを拡張

本格的なワークフローエンジンが必要になる前の段階、つまり「Excelでは限界だが、Temporal/Airflowを導入するほどでもない」というフェーズにおいて、PocketBaseは実用的な選択肢でした。

この記事をシェアする

関連記事