•13 min read

Cloudflare Workers Durable Objects: Distributed Coordination, WebSockets & Edge State

Cloudflare Workers Durable Objects: Distributed Coordination, WebSockets & Edge State

Cloudflare Durable Objects (DOs) provide strongly consistent, low-latency storage and compute at the edge, enabling the construction of globally distributed, stateful applications. This document details architectural patterns for leveraging DOs to manage real-time state, facilitate distributed coordination, and handle high-volume WebSocket connections efficiently.

Audio Briefing
0:00 / 0:00

Durable Objects Fundamentals

A Durable Object is a single-instance class that lives on a single Cloudflare data center. It's uniquely identified by a DurableObjectId and can maintain state across requests. All requests to a specific DO instance are serialized, guaranteeing strong consistency without explicit locking mechanisms. This serialization is a fundamental primitive for distributed coordination.

Core Concepts

  • Single-Instance Guarantee: For a given DurableObjectId, only one instance of the DO class exists globally at any time.
  • Strong Consistency: All operations within a DO instance are atomic and serialized.
  • Transactional Storage: DOs provide an asynchronous key-value store (state.storage) with transactional semantics. Operations within a single state.storage.transaction() block are atomic.
  • Alarms: A mechanism to schedule a one-shot or recurring callback within a DO, enabling background processing.
  • WebSocket Hibernation: DOs can manage thousands of idle WebSocket connections with minimal resource consumption.
Advertisement

Architecture: Real-time Distributed Chat Service

We will construct a simplified, real-time distributed chat service to illustrate these concepts. Each chat room will be a Durable Object. Users connect to a room via WebSockets.

Project Structure

.
├── src
│   ├── index.ts          # Worker entry point, routes requests to DOs
│   ├── chatRoom.ts       # Durable Object for a chat room
│   └── types.ts          # Shared types
├── wrangler.toml         # Cloudflare Workers configuration
└── package.json

wrangler.toml Configuration

Define the Durable Object binding.

name = "chat-service"
main = "src/index.ts"
compatibility_date = "2024-01-01"

[[durable_objects.bindings]]
name = "CHAT_ROOM"
class_name = "ChatRoom"
script_name = "chat-service" # Optional: if DO is in a different worker

src/types.ts

Shared types for messages.

export interface ChatMessage {
  sender: string;
  message: string;
  timestamp: number;
}

export interface WebSocketMessage {
  type: 'chat' | 'join' | 'leave';
  payload: ChatMessage | { user: string };
}

src/chatRoom.ts: Durable Object Implementation

This DO manages a chat room's state, handles WebSocket connections, and persists recent messages.

import { ChatMessage, WebSocketMessage } from './types';

// Maximum number of messages to store in memory and persist
const MAX_MESSAGES = 100;
// Key for storing messages in Durable Object storage
const MESSAGES_STORAGE_KEY = 'messages';

export class ChatRoom {
  private state: DurableObjectState;
  private env: Env; // Environment variables, if any
  private sessions: WebSocket[] = []; // Active WebSocket connections
  private messages: ChatMessage[] = []; // In-memory message buffer

  constructor(state: DurableObjectState, env: Env) {
    this.state = state;
    this.env = env;

    // Restore messages from storage on DO activation
    this.state.blockConcurrencyWhile(async () => {
      const storedMessages = await this.state.storage.get<ChatMessage[]>(MESSAGES_STORAGE_KEY);
      if (storedMessages) {
        this.messages = storedMessages;
      }
      console.log(`ChatRoom ${this.state.id.toString()} initialized with ${this.messages.length} messages.`);
    });
  }

  // Handles incoming HTTP requests, primarily WebSocket upgrades
  async fetch(request: Request): Promise<Response> {
    const url = new URL(request.url);

    if (url.pathname === '/websocket') {
      // Upgrade the request to a WebSocket connection
      const { 0: client, 1: server } = new WebSocketPair();

      this.handleWebSocket(server);

      return new Response(null, { status: 101, webSocket: client });
    } else if (url.pathname === '/messages') {
      // Return recent messages via HTTP for initial load or REST clients
      return new Response(JSON.stringify(this.messages), {
        headers: { 'Content-Type': 'application/json' },
      });
    }

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

  // Handles the lifecycle of a single WebSocket connection
  private async handleWebSocket(server: WebSocket) {
    server.accept();
    this.sessions.push(server);

    // Send historical messages to the new client
    server.send(JSON.stringify({ type: 'history', payload: this.messages }));

    server.addEventListener('message', async (event) => {
      try {
        const message: WebSocketMessage = JSON.parse(event.data as string);
        if (message.type === 'chat') {
          const chatMessage: ChatMessage = {
            ...message.payload as ChatMessage,
            timestamp: Date.now(),
          };
          this.broadcastMessage(chatMessage);
          await this.storeMessage(chatMessage);
        } else if (message.type === 'join') {
          const user = (message.payload as { user: string }).user;
          this.broadcastSystemMessage(`${user} joined the room.`);
        }
      } catch (err) {
        console.error('WebSocket message parsing error:', err);
        server.send(JSON.stringify({ type: 'error', payload: 'Invalid message format.' }));
      }
    });

    server.addEventListener('close', async (event) => {
      console.log(`WebSocket closed: ${event.code} ${event.reason}`);
      this.sessions = this.sessions.filter(s => s !== server);
      // Optionally broadcast leave message
    });

    server.addEventListener('error', (err) => {
      console.error('WebSocket error:', err);
      server.close(1011, 'Internal Error');
    });
  }

  // Broadcasts a message to all connected clients
  private broadcastMessage(message: ChatMessage) {
    const wsMessage: WebSocketMessage = { type: 'chat', payload: message };
    const serializedMessage = JSON.stringify(wsMessage);
    this.sessions.forEach(session => {
      try {
        session.send(serializedMessage);
      } catch (err) {
        console.error('Failed to send message to session:', err);
        // Consider removing broken sessions here
      }
    });
  }

  // Broadcasts a system message (e.g., join/leave)
  private broadcastSystemMessage(text: string) {
    const systemMessage: ChatMessage = { sender: 'System', message: text, timestamp: Date.now() };
    const wsMessage: WebSocketMessage = { type: 'chat', payload: systemMessage };
    const serializedMessage = JSON.stringify(wsMessage);
    this.sessions.forEach(session => {
      try {
        session.send(serializedMessage);
      } catch (err) {
        console.error('Failed to send system message to session:', err);
      }
    });
  }

  // Stores a message in memory and persists it to Durable Object storage
  private async storeMessage(message: ChatMessage) {
    this.messages.push(message);
    if (this.messages.length > MAX_MESSAGES) {
      this.messages.shift(); // Keep only the latest MAX_MESSAGES
    }

    // Persist messages transactionally
    await this.state.storage.transaction(async (txn) => {
      txn.put(MESSAGES_STORAGE_KEY, this.messages);
    });

    // Set an alarm to prune old messages or perform other background tasks
    // This demonstrates using alarms for cron-like functionality
    const currentAlarm = await this.state.storage.getAlarm();
    if (currentAlarm === null || currentAlarm < Date.now() + 60 * 1000) { // Set alarm for 1 minute from now
      await this.state.storage.setAlarm(Date.now() + 60 * 1000);
    }
  }

  // Alarm handler for background tasks
  async alarm() {
    console.log(`Alarm triggered for ChatRoom ${this.state.id.toString()}. Pruning messages.`);
    // Example: Prune messages older than 24 hours
    const twentyFourHoursAgo = Date.now() - (24 * 60 * 60 * 1000);
    const initialLength = this.messages.length;
    this.messages = this.messages.filter(msg => msg.timestamp > twentyFourHoursAgo);

    if (this.messages.length < initialLength) {
      // Only persist if messages were actually pruned
      await this.state.storage.transaction(async (txn) => {
        txn.put(MESSAGES_STORAGE_KEY, this.messages);
      });
      console.log(`Pruned ${initialLength - this.messages.length} messages.`);
    }

    // Optionally reschedule the alarm for the next pruning cycle
    await this.state.storage.setAlarm(Date.now() + 24 * 60 * 60 * 1000); // Reschedule for 24 hours later
  }
}

// Define the environment interface for type safety
interface Env {
  CHAT_ROOM: DurableObjectNamespace;
}

src/index.ts: Worker Entry Point

The Worker acts as a router, dispatching requests to the appropriate Durable Object instance.

import { ChatRoom } from './chatRoom';

interface Env {
  CHAT_ROOM: DurableObjectNamespace;
}

export default {
  async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
    const url = new URL(request.url);
    const roomId = url.pathname.split('/')[1] || 'default-room'; // Extract room ID from path

    // Get a Durable Object ID for the room
    const id = env.CHAT_ROOM.idFromName(roomId);

    // Get the Durable Object stub
    const stub = env.CHAT_ROOM.get(id);

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

// Export the Durable Object class for Wrangler to discover
export { ChatRoom };

Key Architectural Patterns & Features

1. WebSocket Hibernation

Durable Objects excel at managing a large number of idle WebSocket connections. When a WebSocket is idle (no messages sent or received), Cloudflare automatically "hibernates" it. The DO instance remains active, but the CPU cost associated with the connection is effectively zero. Only when a message is sent or received does the connection "wake up," incurring CPU usage. This allows a single DO instance to manage tens of thousands of concurrent, idle connections.

Implementation in ChatRoom: The sessions: WebSocket[] array holds references to active WebSocket connections. Cloudflare's runtime handles the hibernation transparently. The addEventListener callbacks are only invoked when activity occurs.

2. Transactional Storage (SQLite-like Semantics)

The state.storage API provides a key-value store with strong consistency. The transaction() method allows for atomic, all-or-nothing updates. This is crucial for maintaining data integrity.

Implementation in ChatRoom:

await this.state.storage.transaction(async (txn) => {
  txn.put(MESSAGES_STORAGE_KEY, this.messages);
});

This ensures that this.messages is either fully updated in storage or not at all, preventing partial writes if the DO crashes or an error occurs during the transaction. While not a full relational database, transaction() provides a powerful primitive for consistent state management.

3. Alarms for Background Processing

Alarms enable scheduled, one-shot or recurring tasks within a Durable Object. This is ideal for cron-like jobs, data cleanup, or periodic aggregations.

Implementation in ChatRoom:

// Setting an alarm
await this.state.storage.setAlarm(Date.now() + 60 * 1000); // 1 minute from now

// Alarm handler
async alarm() {
  // ... background task logic ...
  await this.state.storage.setAlarm(Date.now() + 24 * 60 * 60 * 1000); // Reschedule
}

The alarm() method is invoked by the Cloudflare runtime when the scheduled time arrives. It's guaranteed to run at most once. If the DO is evicted, the alarm persists and will trigger when the DO is reactivated.

4. Distributed Coordination

The single-instance guarantee of a Durable Object is its most powerful feature for distributed coordination. Any logic requiring a global lock, a single source of truth, or serialized access to a resource can be encapsulated within a DO.

Example: A leader election mechanism. A LeaderElectionDO could manage which worker is the current leader, and other workers would query this DO to determine the leader. The DO's internal state would be the authoritative source.

Example: Rate limiting. A RateLimiterDO could track API call counts for a given user or IP, ensuring global limits are enforced across all edge locations.

Tradeoffs & Considerations

FeatureDurable Objects (DO)Traditional Backend (e.g., EC2 + Redis)Serverless Functions (e.g., Lambda)
ConsistencyStrong (single instance, serialized requests)Eventual (Redis), Strong (DB with locks)Eventual (if shared state is external)
LatencyExtremely low (edge compute)Variable (depends on region, network hops)Variable (cold starts, network hops)
State ManagementBuilt-in transactional KV store, in-memoryExternal (Redis, DB), in-memory (per instance)External (DB, S3, DynamoDB)
ScalingAutomatic per-object, horizontal scaling of objectsManual/Auto-scaling groups, complex state replicationAutomatic per-request, stateless by design
WebSocket SupportFirst-class, hibernation for idle connectionsRequires dedicated servers, complex load balancingLimited, often requires API Gateway + external service
Cost ModelRequest-based, storage, CPU time (incl. hibernation)Instance hours, data transfer, managed service feesRequest-based, duration, memory
ComplexitySimpler for stateful edge logic, less infra managementHigh (infra, ops, scaling, HA)Moderate (integrations, state management)
Use CasesReal-time apps, gaming, chat, distributed locks, CRDTsTraditional web apps, complex databasesEvent-driven, batch processing, APIs
Advertisement

Production Gotchas & Troubleshooting

  1. DO Eviction and Cold Starts: While DOs are "durable," the underlying VM can be evicted and re-instantiated. This means your DO's constructor will run again. Ensure your state.blockConcurrencyWhile() logic is idempotent and efficiently reloads necessary state from state.storage. If your DO holds large in-memory state not persisted, it will be lost on eviction.

    • Fix: Persist critical state to state.storage. Use state.blockConcurrencyWhile to load state synchronously during construction.
    • Symptom: Intermittent data loss or slow initial responses after periods of inactivity.
  2. state.storage.transaction() Failures: Transactions can fail due to contention or internal errors. The async (txn) => { ... } callback will be retried by the DO runtime.

    • Fix: Ensure your transaction logic is idempotent and handles retries gracefully. Avoid side effects within the transaction callback that shouldn't be re-executed.
    • Symptom: TransactionAborted errors in logs, or inconsistent state if not handled correctly.
  3. WebSocket Disconnections: Clients will disconnect. Your DO must handle these gracefully by removing the session from its active list. Unhandled disconnections can lead to memory leaks (holding onto closed WebSocket objects) or attempts to send messages to dead connections.

    • Fix: Implement server.addEventListener('close', ...) and server.addEventListener('error', ...) to clean up sessions.
    • Symptom: Failed to send message to session errors, increasing memory usage over time, or messages not reaching all active clients.
  4. Alarm Scheduling: If an alarm is set for a time in the past, it will trigger immediately. If you set an alarm and then immediately set another one, the second one might overwrite the first, or both might trigger depending on the exact timing and runtime behavior.

    • Fix: Always check await this.state.storage.getAlarm() before setting a new alarm to avoid redundant or conflicting alarms, especially for recurring tasks. Ensure your alarm logic is robust to potential re-triggers.
    • Symptom: Alarms firing unexpectedly or not firing at all.
  5. CPU Usage and Cost: While WebSocket hibernation is efficient for idle connections, active connections and heavy computation within a DO will consume CPU. A single DO instance has finite resources. If a single DO becomes a hotspot for too much active traffic or complex computation, it can become a bottleneck.

    • Fix: Design your application to shard state across multiple DOs (e.g., one DO per chat room, one DO per user). Profile and optimize expensive operations. Consider offloading heavy computation to other Workers or external services if it doesn't require strong consistency within the DO.
    • Symptom: High CPU usage reported in Cloudflare dashboard, increased latency for requests to that DO, or DurableObjectStorageError: Too many concurrent operations errors.

Frequently Asked Questions

  1. Can Durable Objects communicate with each other? Yes. A Durable Object can obtain a stub for another Durable Object (even of a different class) using env.OTHER_DO_NAMESPACE.idFromName(id).get(id) and then fetch() it. This allows for complex inter-DO communication and coordination patterns.

  2. How do I handle schema migrations for state.storage? state.storage is a key-value store, so schema migrations are manual. When you change the structure of data stored under a key, you'll need to implement logic in your DO's constructor or fetch method to detect the old schema and migrate it to the new one. This often involves reading the old data, transforming it, writing the new data, and optionally deleting the old key.

  3. What are the limits on Durable Object storage? Each Durable Object can store up to 128KB of data in state.storage. This limit applies to the total size of all key-value pairs. For larger data, consider storing references to R2 objects or other external storage, and using the DO for metadata and coordination.

  4. How do I ensure my Durable Object is always available? Durable Objects are inherently highly available. Cloudflare's runtime automatically migrates DOs between data centers and restarts them if the underlying machine fails. You don't need to implement explicit failover logic. The state.blockConcurrencyWhile() in the constructor is crucial for ensuring the DO is ready to serve requests after any activation or migration.

  5. Can I use Durable Objects for long-running computations? While DOs can run background tasks via alarms, they are not designed for extremely long-running, CPU-intensive computations. The primary purpose is state management and coordination. If a computation takes too long, it might hit Cloudflare's execution limits or incur high CPU costs. For heavy batch processing, consider offloading to a dedicated compute service.

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