•19 min read

CloudflareWorkers、KV、DurableObjectsで本番MCPサーバーを構築する

CloudflareWorkers、KV、DurableObjectsで本番MCPサーバーを構築する

このガイドでは、Cloudflareのエッジコンピューティングプラットフォームを活用したModel Context Protocol(MCP)サーバーの構築について詳しく説明します。SSEとJSON-RPC 2.0ストリーミングトランスポートを実装し、Durable Objectsでエージェントの状態を管理し、Workers KVにツール結果をキャッシュし、Cloudflare AccessとmTLSを使用してエンドポイントを保護します。

Audio Briefing
0:00 / 0:00

Cloudflare上でのMCPサーバーアーキテクチャ

Cloudflare Workers上のMCPサーバーアーキテクチャは、低遅延で高スループットのエージェントとモデル間の通信のために設計されています。これは以下で構成されます。

  1. Cloudflare Worker: リクエストルーティング、認証、トランスポート層のネゴシエーションを処理する主要なエントリーポイントです。ステートフルな操作のためにDurable Objectsへのプロキシとして機能し、キャッシュのためにKVと連携します。
  2. Cloudflare Durable Objects: 各エージェントセッションに対して、強力な一貫性を持つ、グローバルに一意でステートフルなシングルトンを提供します。これは、会話のコンテキストを維持し、ツール実行のライフサイクルを管理するために不可欠です。
  3. Cloudflare KV: 冪等なツール実行結果や頻繁にアクセスされるエージェント設定をキャッシュするために使用されるキーバリューストアです。
  4. Cloudflare Access & mTLS: Workerエンドポイントを保護し、認証されたエージェントまたはクライアントのみが接続および通信できるようにします。

トランスポート層の実装

MCPは、ストリーミングとリクエスト/レスポンスの両方のパターンを指定しています。ここでは、エージェントのレスポンスをストリーミングするためのSSEと、ストリーミングを含む構造化されたリクエストとレスポンスのためのHTTP上のJSON-RPC 2.0に焦点を当てます。

Server-Sent Events (SSE)

SSEは、サーバーからクライアントへの単方向の永続的な接続を提供し、モデルの出力をストリーミングするのに理想的です。

// src/worker.ts
import { DurableObjectNamespace } from '@cloudflare/workers-types';

interface Env {
  AGENT_SESSION: DurableObjectNamespace;
  TOOL_CACHE: KVNamespace;
  // Cloudflare Access configuration, e.g., AUDIENCE_TAG
  AUDIENCE_TAG: string;
}

export default {
  async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
    const url = new URL(request.url);

    // mTLS and Cloudflare Access validation
    const clientCert = request.headers.get('CF-Client-Cert');
    if (!clientCert) {
      return new Response('mTLS certificate required', { status: 401 });
    }
    // Further validation of clientCert against expected CAs or subject DNs
    // For Cloudflare Access, validate JWT from CF-Access-Jwt-Assertion header
    const jwt = request.headers.get('CF-Access-Jwt-Assertion');
    if (!jwt || !await validateCloudflareAccessJWT(jwt, env.AUDIENCE_TAG)) {
      return new Response('Unauthorized', { status: 401 });
    }

    if (url.pathname.startsWith('/agent/stream/')) {
      const agentId = url.pathname.split('/')[3];
      if (!agentId) {
        return new Response('Agent ID required', { status: 400 });
      }

      // Get Durable Object for the agent session
      const id = env.AGENT_SESSION.idFromName(agentId);
      const stub = env.AGENT_SESSION.get(id);

      // Forward the request to the Durable Object
      return stub.fetch(request);
    }

    // Handle JSON-RPC 2.0 or other requests
    if (url.pathname.startsWith('/agent/rpc/')) {
      const agentId = url.pathname.split('/')[3];
      if (!agentId) {
        return new Response('Agent ID required', { status: 400 });
      }

      const id = env.AGENT_SESSION.idFromName(agentId);
      const stub = env.AGENT_SESSION.get(id);
      return stub.fetch(request);
    }

    return new Response('Not Found', { status: 404 });
  },
};

// Placeholder for Cloudflare Access JWT validation
// In a production environment, use a library or Cloudflare's own validation endpoint
async function validateCloudflareAccessJWT(jwt: string, audienceTag: string): Promise<boolean> {
  // Example: Fetch JWKS from Cloudflare Access and verify signature and audience
  // This is a simplified example. Production code requires robust JWT validation.
  try {
    const response = await fetch(`https://your-team.cloudflareaccess.com/cdn-cgi/access/certs`);
    const jwks = await response.json();
    // ... perform JWT validation using a library like 'jose' ...
    // Check 'aud' claim against audienceTag
    return true; // Placeholder
  } catch (e) {
    console.error('JWT validation failed:', e);
    return false;
  }
}

JSON-RPC 2.0ストリーミング

JSON-RPC 2.0は、単一のHTTP接続を介して複数のJSON-RPC通知またはレスポンスを送信することでストリーミングに適応できます。これは通常、Transfer-Encoding: chunkedを使用するか、SSEのdata:フィールドを活用することで行われます。シンプルさのため、また既存のSSEインフラストラクチャを活用するため、Durable ObjectがSSEを介してJSON-RPCメッセージをプッシュする方法を示します。

// src/durable_object.ts
import { DurableObjectState, DurableObjectEnv } from '@cloudflare/workers-types';

interface AgentSessionState {
  // Example: current model context, tool states, etc.
  context: string[];
  toolResults: Record<string, any>;
}

export class AgentSession {
  state: DurableObjectState;
  env: DurableObjectEnv;
  private sessions: WebSocket[] = []; // For WebSocket-based streaming (alternative to SSE)
  private sseControllers: ReadableStreamDefaultController<Uint8Array>[] = []; // For SSE

  constructor(state: DurableObjectState, env: DurableObjectEnv) {
    this.state = state;
    this.env = env;
    this.state.blockConcurrencyWhile(async () => {
      // Initialize state from storage if available
      const storedState = await this.state.storage.get<AgentSessionState>('agentState');
      if (storedState) {
        // Restore state
      } else {
        // Initialize new state
        await this.state.storage.put('agentState', { context: [], toolResults: {} });
      }
    });
  }

  async fetch(request: Request): Promise<Response> {
    const url = new URL(request.url);

    if (url.pathname.endsWith('/stream')) {
      // Handle SSE connection
      const { readable, writable } = new TransformStream();
      const writer = writable.getWriter();

      // SSE headers
      const headers = {
        'Content-Type': 'text/event-stream',
        'Cache-Control': 'no-cache',
        'Connection': 'keep-alive',
      };

      // Create a controller to push data
      const controller = new ReadableStreamDefaultController<Uint8Array>();
      this.sseControllers.push(controller);

      // Send initial connection message
      await writer.write(new TextEncoder().encode('event: connected\ndata: {}\n\n'));

      // Clean up on disconnect
      request.signal.addEventListener('abort', () => {
        const index = this.sseControllers.indexOf(controller);
        if (index > -1) {
          this.sseControllers.splice(index, 1);
        }
        writer.close();
      });

      return new Response(new ReadableStream({
        start(c) {
          Object.assign(controller, c); // Assign the actual controller
        },
        cancel() {
          const index = this.sseControllers.indexOf(controller);
          if (index > -1) {
            this.sseControllers.splice(index, 1);
          }
        }
      }), { headers });

    } else if (url.pathname.endsWith('/rpc')) {
      // Handle JSON-RPC 2.0 requests
      if (request.method !== 'POST') {
        return new Response('Method Not Allowed', { status: 405 });
      }

      try {
        const rpcRequest = await request.json() as { jsonrpc: string; method: string; params?: any; id?: string | number; };

        if (rpcRequest.jsonrpc !== '2.0') {
          return this.jsonRpcErrorResponse(-32600, 'Invalid Request', rpcRequest.id);
        }

        let rpcResponse: any;
        switch (rpcRequest.method) {
          case 'agent.executeTool':
            rpcResponse = await this.executeTool(rpcRequest.params, rpcRequest.id);
            break;
          case 'agent.getContext':
            rpcResponse = await this.getContext(rpcRequest.id);
            break;
          case 'agent.sendMessage':
            rpcResponse = await this.sendMessage(rpcRequest.params, rpcRequest.id);
            break;
          default:
            rpcResponse = this.jsonRpcErrorResponse(-32601, 'Method not found', rpcRequest.id);
        }

        // Push updates to SSE clients if applicable
        if (rpcRequest.method === 'agent.sendMessage') {
          this.broadcastSSE('agentMessage', rpcResponse.result);
        }

        return new Response(JSON.stringify(rpcResponse), {
          headers: { 'Content-Type': 'application/json' },
        });

      } catch (e: any) {
        console.error('JSON-RPC error:', e);
        return this.jsonRpcErrorResponse(-32700, 'Parse error', null);
      }
    }

    return new Response('Not Found', { status: 404 });
  }

  private async executeTool(params: any, id?: string | number): Promise<any> {
    const { toolName, args } = params;
    const cacheKey = `tool:${toolName}:${JSON.stringify(args)}`;

    // Check KV cache first
    const cachedResult = await this.env.TOOL_CACHE.get(cacheKey);
    if (cachedResult) {
      this.broadcastSSE('toolCacheHit', { toolName, args, result: JSON.parse(cachedResult) });
      return { jsonrpc: '2.0', result: JSON.parse(cachedResult), id };
    }

    // Simulate tool execution
    console.log(`Executing tool: ${toolName} with args:`, args);
    await new Promise(resolve => setTimeout(resolve, 500)); // Simulate work

    const result = {
      output: `Result of ${toolName} with args ${JSON.stringify(args)}`,
      timestamp: Date.now(),
    };

    // Store result in KV cache (e.g., for 1 hour)
    await this.env.TOOL_CACHE.put(cacheKey, JSON.stringify(result), { expirationTtl: 3600 });

    // Update Durable Object state
    const agentState = await this.state.storage.get<AgentSessionState>('agentState');
    if (agentState) {
      agentState.toolResults[toolName] = result;
      await this.state.storage.put('agentState', agentState);
    }

    this.broadcastSSE('toolExecuted', { toolName, args, result });
    return { jsonrpc: '2.0', result, id };
  }

  private async getContext(id?: string | number): Promise<any> {
    const agentState = await this.state.storage.get<AgentSessionState>('agentState');
    return { jsonrpc: '2.0', result: agentState?.context || [], id };
  }

  private async sendMessage(params: any, id?: string | number): Promise<any> {
    const { message } = params;
    // Simulate model processing and response
    console.log('Agent received message:', message);
    await new Promise(resolve => setTimeout(resolve, 200));

    const modelResponse = `Echo: ${message}`;

    // Update Durable Object state
    const agentState = await this.state.storage.get<AgentSessionState>('agentState');
    if (agentState) {
      agentState.context.push(`User: ${message}`);
      agentState.context.push(`Model: ${modelResponse}`);
      await this.state.storage.put('agentState', agentState);
    }

    this.broadcastSSE('modelResponse', { message: modelResponse });
    return { jsonrpc: '2.0', result: { response: modelResponse }, id };
  }

  private jsonRpcErrorResponse(code: number, message: string, id: string | number | null): any {
    return {
      jsonrpc: '2.0',
      error: { code, message },
      id,
    };
  }

  private broadcastSSE(event: string, data: any) {
    const encoder = new TextEncoder();
    const message = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`;
    const encodedMessage = encoder.encode(message);

    this.sseControllers.forEach(controller => {
      try {
        controller.enqueue(encodedMessage);
      } catch (e) {
        console.error('Failed to enqueue SSE message:', e);
        // Consider removing disconnected controllers here
      }
    });
  }
}

エージェントセッション管理のためのDurable Objects

各MCPエージェントセッションには永続的な状態が必要です。Durable Objectsは、単一のCloudflareデータセンターに存在するクラスの単一インスタンスを提供することでこれを実現し、特定のIDに対するすべてのリクエストを処理します。これにより、エージェントの状態に対する強力な一貫性が保証されます。

// wrangler.toml
name = "mcp-worker"
main = "src/worker.ts"
compatibility_date = "2024-01-01"
compatibility_flags = ["nodejs_compat"]

[[kv_namespaces]]
binding = "TOOL_CACHE"
id = "YOUR_KV_NAMESPACE_ID" # Replace with your KV namespace ID

[[durable_objects.bindings]]
name = "AGENT_SESSION"
class_name = "AgentSession"
script_name = "mcp-worker" # Refers to the worker itself, where the DO class is defined

Workers KVによるキャッシング

Workers KVは、冪等なツール実行結果をキャッシュするために使用されます。これにより、冗長な計算や外部API呼び出しが削減され、応答時間が改善され、コストが削減されます。

// Inside AgentSession.executeTool method
const cacheKey = `tool:${toolName}:${JSON.stringify(args)}`;
const cachedResult = await this.env.TOOL_CACHE.get(cacheKey);
if (cachedResult) {
  // ... return cached result ...
}
// ... execute tool, then cache result ...
await this.env.TOOL_CACHE.put(cacheKey, JSON.stringify(result), { expirationTtl: 3600 });

Cloudflare AccessとmTLSによるセキュリティ

MCPエンドポイントのセキュリティ確保は最重要です。Cloudflare AccessはID認識型プロキシを提供し、mTLSはクライアントとサーバーの両方がX.509証明書を使用して相互に認証することを保証します。

  1. Cloudflare Access: Worker用にAccess Applicationを設定します。これにより、エンドポイントがIDプロバイダー(IdP)で保護されます。その後、WorkerはCF-Access-Jwt-Assertionヘッダーを検証します。
  2. mTLS: Cloudflareでホスト名に対してmTLSを有効にします。クライアントは、Cloudflareで設定された信頼できるCAによって発行された有効なクライアント証明書を提示する必要があります。Workerは、証明書の詳細についてCF-Client-Certヘッダーを検査できます。
// src/worker.ts (excerpt)
// mTLS and Cloudflare Access validation
const clientCert = request.headers.get('CF-Client-Cert');
if (!clientCert) {
  return new Response('mTLS certificate required', { status: 401 });
}
// Further validation of clientCert against expected CAs or subject DNs
const jwt = request.headers.get('CF-Access-Jwt-Assertion');
if (!jwt || !await validateCloudflareAccessJWT(jwt, env.AUDIENCE_TAG)) {
  return new Response('Unauthorized', { status: 401 });
}
Advertisement

アーキテクチャのトレードオフ

機能Cloudflare Workers/DO従来のバックエンド(例:Node.js/Redis)
スケーラビリティ自動、グローバル、リクエストごと手動スケーリング、ロードバランサー、地域デプロイが必要
状態管理Durable Objects(強力な一貫性、IDごとにシングルスレッド)Redis/Postgres(外部サービス、結果整合性の可能性あり)
レイテンシエッジ実行、低レイテンシ集中型データセンターのため高レイテンシ
コストモデルリクエストごと、計算時間ごと、ストレージごとサーバーインスタンス、帯域幅、データベース、運用
開発者体験TypeScriptファースト、wrangler CLI、統合プラットフォームより広範なエコシステム、OS/ランタイムのより詳細な制御
セキュリティ統合されたAccess/mTLS、DDoS保護WAF、VPN、認証サービスの手動設定が必要
ベンダーロックインDurable Objects、KVでは高い低い、よりポータブルなコード

本番環境での注意点とトラブルシューティング

  1. Durable Objectのコールドスタート: 非アクティブ期間後、Durable Object IDへの最初のリクエストは、DOインスタンスが初期化されるため、レイテンシが高くなる可能性があります。
    • 解決策: 重要なエージェントの場合、定期的にダミーリクエストを送信してDOを「ウォームアップ」することを検討してください。クライアントは初期のレイテンシスパイクを許容するように設計してください。
  2. KVの結果整合性: KVは高速ですが、結果整合性があります。KVに書き込み、すぐに別のWorkerインスタンスから読み取ると、古いデータが取得される可能性があります。
    • 解決策: 強力な一貫性が必要なデータは、Durable Objectのストレージに直接保存してください。KVは、古さが許容されるキャッシュのようなデータに最適です。
  3. SSE接続管理: ブラウザのSSEクライアントは自動的に再接続します。カスタムクライアントには、指数関数的バックオフを備えた堅牢な再接続ロジックが必要です。
    • 解決策: Last-Event-IDを使用してクライアント側の再接続を実装し、最後の既知のイベントからストリームを再開します。Durable Objectは、Last-Event-IDが存在する場合に再生するためにイベントの履歴を保存する必要があります。
  4. Durable ObjectのblockConcurrencyWhileの誤用: blockConcurrencyWhileを過度に使用したり、その中で長時間実行される操作を実行したりすると、デッドロックや同じDOへの後続のリクエストのレイテンシ増加につながる可能性があります。
    • 解決策: blockConcurrencyWhileは重要な初期化または状態復元にのみ使用してください。長時間実行されるタスク(例:外部API呼び出し)はそれ以外で実行するか、並行バックグラウンドタスクにはPromise.allとstate.waitUntilを使用してください。
  5. mTLS証明書のローテーション: クライアント証明書には有効期限があります。
    • 解決策: エージェントの堅牢な証明書ローテーション戦略を実装してください。Cloudflare AccessはクライアントIDの管理を支援し、直接的な証明書管理の負担を軽減できます。
  6. Cloudflare Access JWT検証: JWTの手動検証は複雑です。
    • 解決策: 実績のあるJWTライブラリ(例:jose)を使用し、正しいCloudflare AccessエンドポイントからJWKSを取得し、すべてのクレーム(aud、iss、exp、nbf)を検証するようにしてください。

よくある質問

  1. Workerの制限を超える長時間実行されるツール実行をどのように処理しますか? Durable Objectsは、呼び出しごとに30秒のCPU時間制限がありますが、接続を最大5分間開いたままにすることができます。本当に長時間実行されるタスク(例:数時間)の場合は、それらを外部サービス(例:Cloudflare Queuesのようなキュー処理システムまたは専用のバックエンド)にオフロードし、Durable Objectを使用して結果をポーリングしたり、Webhookを受信したりします。その後、DOはSSEを介して結果をプッシュできます。

  2. ストリーミングにSSEの代わりにWebSocketを使用できますか? はい、Cloudflare WorkersはWebSocketを完全にサポートしています。Durable ObjectへのWebSocket接続を確立し、DOは接続されたWebSocketオブジェクトのリストを管理してメッセージをブロードキャストします。SSEは単方向のサーバーからクライアントへのストリーミングにはよりシンプルですが、WebSocketは双方向の低遅延通信に適しています。

  3. 単一のDurable Objectへの複数の同時リクエストをどのように管理しますか? Durable Objectsは、デフォルトでリクエストを直列に処理し、強力な一貫性を保証します。リクエストにawait呼び出しが含まれる場合、そのawaitまでは他のリクエストを並行して処理できます。CPUバウンドのタスクの場合、それらは真に直列化されます。これにより状態管理は簡素化されますが、リクエストが譲歩するように設計されていない場合、ボトルネックになる可能性があります。DOメソッドを効率的に設計し、ブロッキング操作を避けてください。

  4. 大規模なエージェントコンテキストやモデルの状態を保存する最良の方法は何ですか? Durable Objectのストレージには制限があります(現在、キーあたり128KB、単一トランザクションで合計1MB、単一DOで合計10MB)。より大きな状態の場合は、以下を検討してください。

    • チャンク化: 大きなオブジェクトを複数のキーに分割します。
    • 外部ストレージ: 大きなブロブをR2(CloudflareのS3互換ストレージ)に保存し、R2オブジェクトポインタをDurable Objectに保存します。
    • 圧縮: データを保存する前に圧縮します。
  5. Durable Objectsをローカルでテストするにはどうすればよいですか? wrangler devはDurable Objectsをサポートしており、DOロジックとWorkerとの相互作用をローカルでテストできます。DO環境をシミュレートするため、ローカル開発が効率的になります。

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement