•13 min read

Building Production MCP Servers on Cloudflare Workers, KV & Durable Objects

Building Production MCP Servers on Cloudflare Workers, KV & Durable Objects

This guide details the construction of Model Context Protocol (MCP) servers leveraging Cloudflare's edge computing platform. We will implement SSE and JSON-RPC 2.0 streaming transports, manage agent state with Durable Objects, cache tool results in Workers KV, and secure the endpoint using Cloudflare Access and mTLS.

Audio Briefing
0:00 / 0:00

MCP Server Architecture on Cloudflare

The MCP server architecture on Cloudflare Workers is designed for low-latency, high-throughput agent-to-model communication. It comprises:

  1. Cloudflare Worker: The primary entry point, handling request routing, authentication, and transport layer negotiation. It acts as a proxy to Durable Objects for stateful operations and interacts with KV for caching.
  2. Cloudflare Durable Objects: Provides strongly consistent, globally unique, and stateful singletons for each agent session. This is crucial for maintaining conversational context and managing tool execution lifecycles.
  3. Cloudflare KV: A key-value store used for caching idempotent tool execution results or frequently accessed agent configuration.
  4. Cloudflare Access & mTLS: Secures the Worker endpoint, ensuring only authorized agents or clients can connect and communicate.

Transport Layer Implementation

MCP specifies both streaming and request/response patterns. We will focus on SSE for streaming agent responses and JSON-RPC 2.0 over HTTP for structured requests and responses, including streaming.

Server-Sent Events (SSE)

SSE provides a unidirectional, persistent connection from the server to the client, ideal for streaming model outputs.

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

JSON-RPC 2.0 can be adapted for streaming by sending multiple JSON-RPC notifications or responses over a single HTTP connection, typically using Transfer-Encoding: chunked or by leveraging SSE's data: field. For simplicity and leveraging existing SSE infrastructure, we'll demonstrate how a Durable Object can push JSON-RPC messages via 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 for Agent Session Management

Each MCP agent session requires persistent state. Durable Objects provide this by offering a single instance of a class that lives on a single Cloudflare data center, handling all requests for a given ID. This ensures strong consistency for agent state.

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

Caching with Workers KV

Workers KV is used for caching idempotent tool execution results. This reduces redundant computation and external API calls, improving response times and reducing costs.

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

Security with Cloudflare Access and mTLS

Securing MCP endpoints is paramount. Cloudflare Access provides identity-aware proxying, and mTLS ensures that both the client and server authenticate each other using X.509 certificates.

  1. Cloudflare Access: Configure an Access Application for your Worker. This will protect the endpoint with your identity provider (IdP). The Worker then validates the CF-Access-Jwt-Assertion header.
  2. mTLS: Enable mTLS for your hostname in Cloudflare. Clients must present a valid client certificate issued by a trusted CA configured in Cloudflare. The Worker can inspect the CF-Client-Cert header for certificate details.
// 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

Architecture Tradeoffs

FeatureCloudflare Workers/DOTraditional Backend (e.g., Node.js/Redis)
ScalabilityAutomatic, global, per-requestRequires manual scaling, load balancers, regional deployment
State ManagementDurable Objects (strong consistency, single-threaded per ID)Redis/Postgres (external service, eventual consistency possible)
LatencyEdge execution, low latencyHigher latency due to centralized data centers
Cost ModelPer-request, per-compute-time, per-storageServer instances, bandwidth, database, ops
Developer ExperienceTypeScript-first, wrangler CLI, integrated platformBroader ecosystem, more control over OS/runtime
SecurityIntegrated Access/mTLS, DDoS protectionRequires manual configuration of WAF, VPN, auth services
Vendor Lock-inHigh for Durable Objects, KVLower, more portable code

Production Gotchas & Troubleshooting

  1. Durable Object Cold Starts: The first request to a Durable Object ID after a period of inactivity can experience higher latency as the DO instance is initialized.
    • Fix: For critical agents, consider "warming up" DOs by sending a dummy request periodically. Design clients to tolerate initial latency spikes.
  2. KV Eventual Consistency: While KV is fast, it's eventually consistent. If you write to KV and immediately read from a different Worker instance, you might get stale data.
    • Fix: For data requiring strong consistency, store it directly in the Durable Object's storage. KV is best for cache-like data where staleness is acceptable.
  3. SSE Connection Management: Browser SSE clients automatically reconnect. Custom clients need robust reconnection logic with exponential backoff.
    • Fix: Implement client-side reconnection with Last-Event-ID to resume streams from the last known event. The Durable Object should store a history of events to replay if Last-Event-ID is present.
  4. Durable Object blockConcurrencyWhile Misuse: Overusing blockConcurrencyWhile or performing long-running operations within it can lead to deadlocks or increased latency for subsequent requests to the same DO.
    • Fix: Only use blockConcurrencyWhile for critical initialization or state restoration. Perform long-running tasks (e.g., external API calls) outside of it, or use Promise.all with state.waitUntil for concurrent background tasks.
  5. mTLS Certificate Rotation: Client certificates have expiration dates.
    • Fix: Implement a robust certificate rotation strategy for your agents. Cloudflare Access can help manage client identities, reducing the burden of direct certificate management.
  6. Cloudflare Access JWT Validation: Manually validating JWTs is complex.
    • Fix: Use a battle-tested JWT library (e.g., jose) and ensure you fetch the JWKS from the correct Cloudflare Access endpoint and validate all claims (aud, iss, exp, nbf).

Frequently Asked Questions

  1. How do I handle long-running tool executions that exceed Worker limits? Durable Objects have a 30-second CPU time limit per invocation, but can keep connections open for up to 5 minutes. For truly long-running tasks (e.g., hours), offload them to an external service (e.g., a queue processing system like Cloudflare Queues or a dedicated backend) and use the Durable Object to poll for results or receive webhooks. The DO can then push the result via SSE.

  2. Can I use WebSockets instead of SSE for streaming? Yes, Cloudflare Workers fully support WebSockets. You would establish a WebSocket connection to the Durable Object, and the DO would manage a list of connected WebSocket objects to broadcast messages. SSE is simpler for unidirectional server-to-client streaming, while WebSockets are better for bidirectional, low-latency communication.

  3. How do I manage multiple concurrent requests to a single Durable Object? Durable Objects process requests serially by default, ensuring strong consistency. If a request involves an await call, other requests can be processed concurrently up to that await. For CPU-bound tasks, they are truly serialized. This simplifies state management but can be a bottleneck if requests are not designed to yield. Design your DO methods to be efficient and avoid blocking operations.

  4. What's the best way to store large agent contexts or model states? Durable Object storage has limits (currently 128KB per key, 1MB total for a single transaction, 10MB total for a single DO). For larger states, consider:

    • Chunking: Split large objects into multiple keys.
    • External Storage: Store large blobs in R2 (Cloudflare's S3-compatible storage) and store R2 object pointers in the Durable Object.
    • Compression: Compress data before storing it.
  5. How do I test Durable Objects locally? wrangler dev supports Durable Objects, allowing you to test your DO logic and interactions with the Worker locally. It simulates the DO environment, making local development efficient.

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