•16 min read

Xây dựng máy chủ MCP sản xuất trên Cloudflare Workers, KV & Durable Objects

Xây dựng máy chủ MCP sản xuất trên Cloudflare Workers, KV & Durable Objects

Hướng dẫn này trình bày chi tiết việc xây dựng các máy chủ Giao thức Ngữ cảnh Mô hình (MCP) tận dụng nền tảng điện toán biên của Cloudflare. Chúng ta sẽ triển khai các giao thức truyền tải luồng SSE và JSON-RPC 2.0, quản lý trạng thái tác nhân bằng Durable Objects, lưu trữ kết quả công cụ trong Workers KV và bảo mật điểm cuối bằng Cloudflare Access và mTLS.

Audio Briefing
0:00 / 0:00

Kiến trúc Máy chủ MCP trên Cloudflare

Kiến trúc máy chủ MCP trên Cloudflare Workers được thiết kế để giao tiếp tác nhân-mô hình có độ trễ thấp, thông lượng cao. Nó bao gồm:

  1. Cloudflare Worker: Điểm vào chính, xử lý định tuyến yêu cầu, xác thực và đàm phán lớp truyền tải. Nó hoạt động như một proxy cho Durable Objects đối với các hoạt động có trạng thái và tương tác với KV để lưu trữ.
  2. Cloudflare Durable Objects: Cung cấp các singleton nhất quán mạnh mẽ, duy nhất trên toàn cầu và có trạng thái cho mỗi phiên tác nhân. Điều này rất quan trọng để duy trì ngữ cảnh hội thoại và quản lý vòng đời thực thi công cụ.
  3. Cloudflare KV: Một kho lưu trữ khóa-giá trị được sử dụng để lưu trữ kết quả thực thi công cụ bất biến hoặc cấu hình tác nhân được truy cập thường xuyên.
  4. Cloudflare Access & mTLS: Bảo mật điểm cuối Worker, đảm bảo chỉ các tác nhân hoặc máy khách được ủy quyền mới có thể kết nối và giao tiếp.

Triển khai Lớp Truyền tải

MCP chỉ định cả mẫu luồng và yêu cầu/phản hồi. Chúng ta sẽ tập trung vào SSE để truyền tải luồng phản hồi của tác nhân và JSON-RPC 2.0 qua HTTP cho các yêu cầu và phản hồi có cấu trúc, bao gồm cả luồng.

Server-Sent Events (SSE)

SSE cung cấp một kết nối một chiều, liên tục từ máy chủ đến máy khách, lý tưởng để truyền tải luồng đầu ra của mô hình.

// 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;
  }
}

Truyền tải luồng JSON-RPC 2.0

JSON-RPC 2.0 có thể được điều chỉnh để truyền tải luồng bằng cách gửi nhiều thông báo hoặc phản hồi JSON-RPC qua một kết nối HTTP duy nhất, thường sử dụng Transfer-Encoding: chunked hoặc bằng cách tận dụng trường data: của SSE. Để đơn giản và tận dụng cơ sở hạ tầng SSE hiện có, chúng ta sẽ trình bày cách Durable Object có thể đẩy các thông báo JSON-RPC qua SSE.

// 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 để Quản lý Phiên Tác nhân

Mỗi phiên tác nhân MCP yêu cầu trạng thái liên tục. Durable Objects cung cấp điều này bằng cách cung cấp một thể hiện duy nhất của một lớp sống trên một trung tâm dữ liệu Cloudflare duy nhất, xử lý tất cả các yêu cầu cho một ID nhất định. Điều này đảm bảo tính nhất quán mạnh mẽ cho trạng thái tác nhân.

// 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

Lưu trữ với Workers KV

Workers KV được sử dụng để lưu trữ kết quả thực thi công cụ bất biến. Điều này làm giảm tính toán dư thừa và các cuộc gọi API bên ngoài, cải thiện thời gian phản hồi và giảm chi phí.

// 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 });

Bảo mật với Cloudflare Access và mTLS

Bảo mật các điểm cuối MCP là tối quan trọng. Cloudflare Access cung cấp ủy quyền nhận dạng, và mTLS đảm bảo rằng cả máy khách và máy chủ đều xác thực lẫn nhau bằng cách sử dụng chứng chỉ X.509.

  1. Cloudflare Access: Cấu hình một Ứng dụng Access cho Worker của bạn. Điều này sẽ bảo vệ điểm cuối bằng nhà cung cấp danh tính (IdP) của bạn. Worker sau đó xác thực tiêu đề CF-Access-Jwt-Assertion.
  2. mTLS: Bật mTLS cho tên máy chủ của bạn trong Cloudflare. Máy khách phải xuất trình chứng chỉ máy khách hợp lệ được cấp bởi một CA đáng tin cậy được cấu hình trong Cloudflare. Worker có thể kiểm tra tiêu đề CF-Client-Cert để biết chi tiết chứng chỉ.
// 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

Đánh đổi Kiến trúc

Tính năngCloudflare Workers/DOBackend truyền thống (ví dụ: Node.js/Redis)
Khả năng mở rộngTự động, toàn cầu, theo yêu cầuYêu cầu mở rộng thủ công, bộ cân bằng tải, triển khai theo khu vực
Quản lý trạng tháiDurable Objects (nhất quán mạnh mẽ, đơn luồng trên mỗi ID)Redis/Postgres (dịch vụ bên ngoài, có thể có tính nhất quán cuối cùng)
Độ trễThực thi biên, độ trễ thấpĐộ trễ cao hơn do các trung tâm dữ liệu tập trung
Mô hình chi phíTheo yêu cầu, theo thời gian tính toán, theo lưu trữCác phiên bản máy chủ, băng thông, cơ sở dữ liệu, hoạt động
Trải nghiệm nhà phát triểnƯu tiên TypeScript, CLI wrangler, nền tảng tích hợpHệ sinh thái rộng hơn, kiểm soát nhiều hơn đối với HĐH/thời gian chạy
Bảo mậtAccess/mTLS tích hợp, bảo vệ DDoSYêu cầu cấu hình thủ công WAF, VPN, dịch vụ xác thực
Khóa nhà cung cấpCao đối với Durable Objects, KVThấp hơn, mã dễ di chuyển hơn

Những vấn đề và cách khắc phục trong sản xuất

  1. Khởi động lạnh Durable Object: Yêu cầu đầu tiên đến một ID Durable Object sau một thời gian không hoạt động có thể gặp độ trễ cao hơn khi thể hiện DO được khởi tạo.
    • Cách khắc phục: Đối với các tác nhân quan trọng, hãy cân nhắc "làm nóng" DO bằng cách gửi một yêu cầu giả định định kỳ. Thiết kế máy khách để chịu được các đợt tăng độ trễ ban đầu.
  2. Tính nhất quán cuối cùng của KV: Mặc dù KV nhanh, nhưng nó có tính nhất quán cuối cùng. Nếu bạn ghi vào KV và đọc ngay lập tức từ một thể hiện Worker khác, bạn có thể nhận được dữ liệu cũ.
    • Cách khắc phục: Đối với dữ liệu yêu cầu tính nhất quán mạnh mẽ, hãy lưu trữ trực tiếp trong bộ nhớ của Durable Object. KV tốt nhất cho dữ liệu giống như bộ nhớ đệm mà độ cũ có thể chấp nhận được.
  3. Quản lý kết nối SSE: Các máy khách SSE của trình duyệt tự động kết nối lại. Các máy khách tùy chỉnh cần logic kết nối lại mạnh mẽ với thời gian chờ tăng dần.
    • Cách khắc phục: Triển khai kết nối lại phía máy khách với Last-Event-ID để tiếp tục luồng từ sự kiện cuối cùng được biết. Durable Object nên lưu trữ lịch sử các sự kiện để phát lại nếu Last-Event-ID có mặt.
  4. Lạm dụng blockConcurrencyWhile của Durable Object: Lạm dụng blockConcurrencyWhile hoặc thực hiện các hoạt động chạy dài trong đó có thể dẫn đến tắc nghẽn hoặc tăng độ trễ cho các yêu cầu tiếp theo đến cùng một DO.
    • Cách khắc phục: Chỉ sử dụng blockConcurrencyWhile cho việc khởi tạo hoặc khôi phục trạng thái quan trọng. Thực hiện các tác vụ chạy dài (ví dụ: gọi API bên ngoài) bên ngoài nó, hoặc sử dụng Promise.all với state.waitUntil cho các tác vụ nền đồng thời.
  5. Xoay vòng chứng chỉ mTLS: Chứng chỉ máy khách có ngày hết hạn.
    • Cách khắc phục: Triển khai chiến lược xoay vòng chứng chỉ mạnh mẽ cho các tác nhân của bạn. Cloudflare Access có thể giúp quản lý danh tính máy khách, giảm gánh nặng quản lý chứng chỉ trực tiếp.
  6. Xác thực JWT của Cloudflare Access: Xác thực JWT thủ công rất phức tạp.
    • Cách khắc phục: Sử dụng thư viện JWT đã được kiểm chứng (ví dụ: jose) và đảm bảo bạn lấy JWKS từ điểm cuối Cloudflare Access chính xác và xác thực tất cả các yêu cầu (aud, iss, exp, nbf).

Các câu hỏi thường gặp

  1. Làm cách nào để xử lý các tác vụ thực thi công cụ chạy dài vượt quá giới hạn của Worker? Durable Objects có giới hạn thời gian CPU là 30 giây cho mỗi lần gọi, nhưng có thể giữ kết nối mở trong tối đa 5 phút. Đối với các tác vụ thực sự chạy dài (ví dụ: hàng giờ), hãy chuyển chúng sang một dịch vụ bên ngoài (ví dụ: hệ thống xử lý hàng đợi như Cloudflare Queues hoặc một backend chuyên dụng) và sử dụng Durable Object để thăm dò kết quả hoặc nhận webhook. DO sau đó có thể đẩy kết quả qua SSE.

  2. Tôi có thể sử dụng WebSockets thay vì SSE để truyền tải luồng không? Có, Cloudflare Workers hỗ trợ đầy đủ WebSockets. Bạn sẽ thiết lập một kết nối WebSocket đến Durable Object, và DO sẽ quản lý danh sách các đối tượng WebSocket được kết nối để phát các thông báo. SSE đơn giản hơn cho việc truyền tải luồng một chiều từ máy chủ đến máy khách, trong khi WebSockets tốt hơn cho giao tiếp hai chiều, độ trễ thấp.

  3. Làm cách nào để quản lý nhiều yêu cầu đồng thời đến một Durable Object duy nhất? Durable Objects xử lý các yêu cầu nối tiếp theo mặc định, đảm bảo tính nhất quán mạnh mẽ. Nếu một yêu cầu liên quan đến một cuộc gọi await, các yêu cầu khác có thể được xử lý đồng thời cho đến await đó. Đối với các tác vụ bị giới hạn bởi CPU, chúng thực sự được nối tiếp. Điều này đơn giản hóa việc quản lý trạng thái nhưng có thể là một nút thắt cổ chai nếu các yêu cầu không được thiết kế để nhường quyền. Thiết kế các phương thức DO của bạn hiệu quả và tránh các hoạt động chặn.

  4. Cách tốt nhất để lưu trữ các ngữ cảnh tác nhân lớn hoặc trạng thái mô hình là gì? Bộ nhớ Durable Object có giới hạn (hiện tại là 128KB mỗi khóa, tổng cộng 1MB cho một giao dịch duy nhất, tổng cộng 10MB cho một DO duy nhất). Đối với các trạng thái lớn hơn, hãy cân nhắc:

    • Phân đoạn: Chia các đối tượng lớn thành nhiều khóa.
    • Bộ nhớ ngoài: Lưu trữ các blob lớn trong R2 (bộ nhớ tương thích S3 của Cloudflare) và lưu trữ các con trỏ đối tượng R2 trong Durable Object.
    • Nén: Nén dữ liệu trước khi lưu trữ.
  5. Làm cách nào để kiểm tra Durable Objects cục bộ? wrangler dev hỗ trợ Durable Objects, cho phép bạn kiểm tra logic DO và tương tác với Worker cục bộ. Nó mô phỏng môi trường DO, giúp phát triển cục bộ hiệu quả.

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