Cloudflare Workflows で Human-in-the-loop を実装する

Cloudflare Workflows で Human-in-the-loop を実装する

以下の記事で少し触れましたが、Cloudflare Workflows を使って、Human-in-the-loop の処理を実装しています。

AI 関連の処理をバッチで動かす場合はリージョンに注意 – もばらぶエンジニアブログ

Human-in-the-loop とは、AI の処理の途中に人間の判断を挟む設計パターンの事です。例えば、AI が受信したメールを読んで返信メールの下書きを書いて、人間に送信して良いか聞いてきて、承認したらそれを送信する、みたいなものです。

Cloudflare Workflows については以下を参照してください。

Overview · Cloudflare Workflows docs

やりたい事

まずは簡単な事から試してみようと思い、以下のような処理を行う事にしました。

  1. Claude Code がネット上から、とあるトピックのニュースを収集
  2. Claude Code が MCP サーバー経由で、Slack の管理者専用チャンネルに「この内容を投稿して良いか」と投稿
  3. 人間が OK/NG/修正依頼の返信
  4. OK の場合は、Claude Code が Slack の一般ユーザー向けチャンネルにニュースを投稿

構成

AI に描いてもらったテキスト形式の図を少し修正して載せます。

Cloudflare Cron Trigger
        │
        ▼
   Worker (/trigger, /scheduled)
        │ Workflow起動
        ▼
  ApprovalWorkflow (Cloudflare Workflows)
    1. propose  … Container(claude -p, WebSearch)でWeb検索・判断・要約
                  D1 proposals の直近30日分の履歴をプロンプトに渡し、重複提案を避ける
                  (提案後に記録し、publish/却下/タイムアウト/修正依頼でstatusを更新)
    2. notify   … Container(claude -p, Slack MCP: Bloque gateway経由)で
                  SLACK_ADMIN_CHANNEL(admin限定のプライベートチャンネル)に提案投稿
    3. D1に thread_ts <-> workflow_id を保存
    4. waitForEvent … admin限定チャンネル側スレッドへの人間の返信を待つ
                      (ここでは課金されない)

~~~~~~~~~~ここが処理の区切り、人間の判断がここで行われる

Slackスレッド返信 (admin限定チャンネル側)
        │ Slack 側で Event Subscriptions の設定
        ▼
Worker (/slack/events) … 署名検証 → D1でworkflow_id検索 → instance.sendEvent()
        │ Workflow起動
        ▼
  ApprovalWorkflow (Cloudflare Workflows)
    5. publish  … 承認されたらContainer(claude -p, Slack MCP: Bloque gateway経由)で
                  SLACK_PUBLIC_CHANNEL(公開チャンネル)に確定版を新規投稿し、
                  admin限定チャンネル側スレッドに完了報告を返信

使っているものと用途を簡単にまとめます。あと、私は Cloudflare は初心者ですが AWS には慣れているので、AWS だと何にあたるのかも書いておきます。

  • Cron
    • 定時起動のトリガー
    • AWS の EventBridge Scheduler 相当
  • Worker
    • 起動経路は以下の3つ
      • Cron: 後述の Workflow の propose を呼び出す
      • HTTP エンドポイント /trigger : 同じく propose を呼び出す
      • HTTP エンドポイント /slack/events : Workflow の publish を呼び出す
    • API Gateway + Lambda 相当?
  • Workflow
    • 処理の本体
    • 上述の1〜4と5という2つの処理の固まりに分かれる
    • AWS だと Bedrock で出来そうですが、調べてません
  • D1
    • テーブルは以下の2つ
      • Workflow の1〜4と5の間でデータをやり取りするためのテーブル
      • 過去の提案を保存しておくテーブル、重複提案を避けるために導入
    • SQLite ベースらしい
  • Container
    • Claude Code をインストールした Docker コンテナ
    • Workflow の途中(ステップ1, 2, 5)で呼び出される
    • Docker イメージを使った Lambda か素の ECS task に相当

Claude Code からは Bloque 経由で Slack MCP サーバーに繋いでいます。

Bloque – All Your MCPs, Centrally Managed.

現在は MCP サーバーが1つだけなのでわざわざ MCP ゲートウェイ(Bloque)を挟む必要はそこまでないのですが、今後、他の MCP サーバーを追加してもっと複雑なフローにするときなどを考慮して、入れてみました。

プログラム・設定ファイル等

気が向けばサンプルプロジェクトとして公開しようと思いますが、とりあえずはプログラム・設定ファイルの重要そうなところを抜粋します。

wrangler.jsonc はほぼそのまま載せます。

{
	"$schema": "node_modules/wrangler/config-schema.json",
	"name": "foo-news-poster",
	"main": "src/index.ts",
	"compatibility_date": "2026-09-01",
	"compatibility_flags": [
		"nodejs_compat"
	],
	// Workers Logs(ダッシュボードで after-the-fact に確認できる永続ログ)。
	// wrangler tail はライブのみ・接続中しか見えないため、非同期に進む
	// Workflow(Slack承認待ちで数時間空くこともある)のデバッグ用に有効化。
	"observability": {
		"enabled": true
	},
	// 定期実行(1日3回の例)。Cloudflare Cron Trigger を使う(GitHub Actions cronは不採用)。
	"triggers": {
		"crons": [
			"0 1,7,13 * * *"
		]
	},
	// human-in-the-loop の中核。承認待ち状態を durable に保持する。
	"workflows": [
		{
			"name": "foo-news-poster-approval",
			"binding": "APPROVAL_WORKFLOW",
			"class_name": "ApprovalWorkflow"
		}
	],
	// Claude Code (`claude -p`) を動かす Container。
	"containers": [
		{
			"class_name": "ClaudeCodeContainer",
			"image": "./container/Dockerfile",
			"image_build_context": ".",
			"instance_type": "basic",
			// Anthropic APIは香港(HKG)等からのリクエストを 403 "Request not allowed" で拒否する。
			// cron起動のWorkflowではコンテナがHKGに配置され `claude -p` が失敗したため、
			// 対応リージョン(北米)に固定する。
			"constraints": {
				"regions": ["ENAM", "WNAM"]
			},
			"max_instances": 5
		}
	],
	"durable_objects": {
		"bindings": [
			{
				"name": "CLAUDE_CONTAINER",
				"class_name": "ClaudeCodeContainer"
			}
		]
	},
	"migrations": [
		{
			"tag": "v1",
			"new_sqlite_classes": [
				"ClaudeCodeContainer"
			]
		}
	],
	// Slack の thread_ts <-> workflow instance の対応管理用(Slack Events受信時の検索に使う)
	"d1_databases": [
		{
			"binding": "DB",
			"database_name": "foo-news-poster-db",
			"database_id": "xxxxxxxx",
			"migrations_dir": "migrations"
		}
	],
	"vars": {
		// 提案(承認待ち)の投稿先。admin限定のプライベートチャンネル/グループ。
		"SLACK_ADMIN_CHANNEL": "C0ADMIN",
		// 承認後、確定版を投稿する公開チャンネル。
		"SLACK_PUBLIC_CHANNEL": "C0PUBLIC"
	}
	// secrets (wrangler secret put で投入。wrangler.jsonc には書かない):
	//   SLACK_SIGNING_SECRET     … Slack Events API の署名検証用
	//   BLOQUE_API_KEY           … Container内のClaude Code(Slack MCP: bloque gateway)が使う
	//   CLAUDE_CODE_OAUTH_TOKEN  … Container内の `claude -p` の認証用(`claude setup-token`で発行、API key の場合は ANTHROPIC_API_KEY)
}

Container の Dockerfile もそのまま載せます。

FROM node:22-slim
# ↑ 24で良いと思います

# Claude Code CLI
RUN npm install -g @anthropic-ai/claude-code

WORKDIR /app

# 実行に必要な最小限のnode依存(型定義とtsxのみ)
COPY container/package.json ./package.json
RUN npm install

COPY container/run.ts ./run.ts
# スキル定義をイメージに同梱(Claude Codeがロードできるように)
COPY .claude ./.claude
# Slack MCP (Bloque gateway経由のslack-bot) の接続設定。
# WORKDIR直下に置くことでClaude Codeが自動でロードする。
COPY .mcp.json ./.mcp.json

ENV PORT=8080
EXPOSE 8080

# CLAUDE_CODE_OAUTH_TOKEN / BLOQUE_API_KEY はコンテナ起動時に環境変数として注入される想定
# (Worker側でContainerを起動する際にenvを渡す。CLAUDE_CODE_OAUTH_TOKENは`claude -p`自体の
#  認証に、BLOQUE_API_KEYは.mcp.json内の ${BLOQUE_API_KEY} 展開に使われ、Bloque gateway
#  経由でSlack MCPに認証する)

CMD ["npx", "tsx", "run.ts"]

Workflow はメイン処理だけほぼそのまま載せます。

  async run(event: WorkflowEvent<WorkflowParams>, step: WorkflowStep) {
    let feedback: string | undefined;

    // 過去の提案履歴(同じニュースを繰り返し提案しないためにproposeへ渡す)。
    // step.do 内で取得するのでリプレイ間で結果が固定される。
    // 今回のrun自身の過去attemptは含めない(修正依頼では同じニュースを扱い得るため)。
    const history = await step.do('load-history', async () => {
      const { results } = await this.env.DB.prepare(
        `SELECT title, source_url, status, created_at FROM proposals
         WHERE created_at >= datetime('now', ?)
         ORDER BY created_at DESC LIMIT ?`,
      )
        .bind(`-${HISTORY_DAYS} days`, HISTORY_LIMIT)
        .all<{ title: string; source_url: string | null; status: ProposalStatus; created_at: string }>();
      return results.map(
        (r): ProposalHistoryItem => ({
          title: r.title,
          sourceUrl: r.source_url ?? undefined,
          status: r.status,
          createdAt: r.created_at,
        }),
      );
    });

    for (let attempt = 0; attempt <= MAX_REVISIONS; attempt++) {
      // --- 1(情報収集) + 2(判断) + 3(要約・提案) ---
      // Claude Code(スキル)がWeb検索 -> 一般人向けかの判断 -> 要約まで行う
      const proposal = await step.do(`propose-${attempt}`, async () => {
        return await this.runClaudeCode<ProposeResult>(
          'propose',
          {
            topic: event.payload.topic,
            feedback, // 修正依頼があれば2回目以降のプロンプトに含める
            history,
          },
          event.instanceId,
        );
      });

      if (proposal.status === 'no_candidate' || !proposal.summary) {
        // 公開に値する一般向け情報が見つからなかった場合はここで終了
        return { result: 'no_candidate' as const };
      }

      // 提案を履歴として記録(UNIQUE + INSERT OR IGNOREでstepリトライ時も冪等)
      await step.do(`record-proposal-${attempt}`, async () => {
        await this.env.DB.prepare(
          `INSERT OR IGNORE INTO proposals (workflow_id, attempt, title, summary, source_url)
           VALUES (?, ?, ?, ?, ?)`,
        )
          .bind(event.instanceId, attempt, proposal.title ?? '', proposal.summary!, proposal.sourceUrl ?? null)
          .run();
      });

      // 提案をSlackに投稿して承認を仰ぐ(Claude CodeがSlack MCP経由で投稿)
      const posted = await step.do(`notify-${attempt}`, async () => {
        return await this.runClaudeCode<NotifyResult>(
          'notify',
          {
            title: proposal.title,
            summary: proposal.summary,
            sourceUrl: proposal.sourceUrl,
          },
          event.instanceId,
        );
      });

      // Slackのスレッド返信を受け取ったときに、どのWorkflowインスタンスを
      // 起こせばよいか検索できるようにD1へ対応関係を保存
      await step.do(`persist-${attempt}`, async () => {
        await this.env.DB.prepare(
          `INSERT OR REPLACE INTO pending_approvals (workflow_id, thread_ts, channel)
           VALUES (?, ?, ?)`,
        )
          .bind(event.instanceId, posted.threadTs, posted.channel)
          .run();
      });

      // --- 4(人間の判断) ---
      // Slack Events APIからのWebhookが `instance.sendEvent()` を呼ぶまでここで待機。
      // 待機中はコンテナも課金も発生しない。
      // 12時間以内に返信がなければタイムアウトし、却下扱いとする。
      let decision: WorkflowStepEvent<Decision>;
      try {
        decision = await step.waitForEvent<Decision>(`await-decision-${attempt}`, {
          type: 'human-decision',
          timeout: '12 hours',
        });
      } catch {
        await this.markProposal(step, event.instanceId, attempt, 'timed_out');
        return { result: 'rejected' as const };
      }

      if (decision.payload.kind === 'approved') {
        // --- 5(処理実行): 正式投稿 ---
        await step.do('publish', async () => {
          return await this.runClaudeCode<PublishResult>(
            'publish',
            {
              title: proposal.title,
              summary: proposal.summary,
              sourceUrl: proposal.sourceUrl,
              // 承認元(SLACK_ADMIN_CHANNEL)のスレッド。公開後にここへ完了報告する。
              adminChannel: posted.channel,
              adminThreadTs: posted.threadTs,
            },
            event.instanceId,
          );
        });
        await this.markProposal(step, event.instanceId, attempt, 'published');
        return { result: 'published' as const };
      }

      if (decision.payload.kind === 'rejected') {
        await this.markProposal(step, event.instanceId, attempt, 'rejected');
        return { result: 'rejected' as const };
      }

      // revise: フィードバックを持って次のループ(再提案)へ
      await this.markProposal(step, event.instanceId, attempt, 'revised');
      feedback = decision.payload.feedback;
    }

    return { result: 'max_revisions_reached' as const };
  }

同様な事が実現できる他の手段

今回の構成は説明した通りですが、HITL を実装する方法は他にも山ほどありますのでいくつか紹介します。

Cloudflare Agents SDK を使って実装する

今回は Cloudflare Workflows を使いましたが、Agents SDK を使うという選択肢もあります。今回書いた処理を少し簡単にしてくれる部分もあったり、web が UI の AI エージェントを実装するのに便利な機能もあるのですが、今回はバッチ処理なのでほぼメリットがないため、採用は見送りました。

ドキュメントは以下にあります。

Agents · Cloudflare Agents docs

LangChain, LangGraph + HumanInTheLoopMiddleware で実装する

LangChain には HumanInTheLoopMiddleware というズバリの名前の機能があります。LangChain から直接使っても良いですし、LangGraph と組み合わせても良いでしょう。

LangChain を使った事がある人・動作を理解している人であれば良い選択肢だと思います。ただ、今回の方法に比べると(最近は AI が実装する事が大半かとは思いますが)実装量は少し多くなります。また、現在の処理状態を保存しておくデータストアも別途用意する必要があります(他の多くの方法でも同様です)。

ドキュメントは以下の通りです。

Human-in-the-loop – Docs by LangChain

Vercel Workflows + Vercel AI SDK

一言でまとめると今回やった事の Vercel 版です。私は Vercel は使っておらず料金とかも調べてないですが、Vercel に慣れている人であればこちらを検討しても良いと思います。

ドキュメントは以下を参照してください。

Vercel Workflows

HITL 実装の cookbook もありました。

Human-in-the-Loop with Next.js | AI SDK

n8n の human review 機能

あまりプログラムを書きたくない人は、n8n 等のオートメーションツール・ノーコードツールを使うという選択肢があります。(私はこの機能を使った事はありませんが)n8n には human review 機能というのがあるので、それで実現できそうです。

以下のページに具体的なやり方が書いてあります。

Human-in-the-loop for tools | Build | n8n Docs

その他にも色々

AI 対応は今熱い分野ですし他にも色々ありますし、これからもどんどん出てくると思います。

感想: 簡単だった

インフラ構築は wrangler.jsonc というファイルに書くだけで良く、プログラムの量も大して多くないので、割と簡単に実装出来ました(以下の問題を除く)。

AI 関連の処理をバッチで動かす場合はリージョンに注意 – もばらぶエンジニアブログ

AWS で実装しようとすると、インフラ構築は CDK だと(メリットはありますが)記述量が多くなってしまいますし、使うサービスも多くなりそうですし、Cloudflare より大がかりな構成となってしまいそうです。

とはいえこれは個人の感想なので、(できれば仕事で)機会があれば、AWS でも似たような事を実装してみたいと思っています。

まとめ

今回初めて Cloudflare を使いましたが、思ったより簡単に HITL が実装出来て驚きました。

今回の内容は正直なところノーコード・ローコードでも実現できたと思いますし、仕事でやるとしたら n8n や Dify なども選択肢になると思います。ただ、ノーコード・ローコードだとちょっとやりにくいような処理を実装したいけど、大がかりな仕組みにはしたくない、という我々のような中小企業にとっては、Cloudflare は良い選択肢だと感じました。

宣伝: 今回のような自動化処理・HITL の実装、AI・MCP 関連でお困りの事があれば、以下のお問い合わせフォームよりお気軽にお問い合わせください。

お問い合わせ – もばらぶん

we are hiring

優秀な技術者と一緒に、好きな場所で働きませんか

株式会社もばらぶでは、優秀で意欲に溢れる方を常に求めています。働く場所は自由、働く時間も柔軟に選択可能です。

現在、以下の職種を募集中です。ご興味のある方は、リンク先をご参照下さい。

コメントを残す