Skip to content
File

Blob: samples/durable-objects-chat/chat.js

javascript548 lines
1// Copyright (c) 2022-2023 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5// This is the Edge Chat Demo Worker, built using Durable Objects!
6 
7// ===============================
8// Introduction to Modules
9// ===============================
10//
11// The first thing you might notice, if you are familiar with the Workers platform, is that this
12// Worker is written differently from others you may have seen. It even has a different file
13// extension. The `mjs` extension means this JavaScript is an ES Module, which, among other things,
14// means it has imports and exports. Unlike other Workers, this code doesn't use
15// `addEventListener("fetch", handler)` to register its main HTTP handler; instead, it _exports_
16// a handler, as we'll see below.
17//
18// This is a new way of writing Workers that we expect to introduce more broadly in the future. We
19// like this syntax because it is *composable*: You can take two workers written this way and
20// merge them into one worker, by importing the two Workers' exported handlers yourself, and then
21// exporting a new handler that call into the other Workers as appropriate.
22//
23// This new syntax is required when using Durable Objects, because your Durable Objects are
24// implemented by classes, and those classes need to be exported. The new syntax can be used for
25// writing regular Workers (without Durable Objects) too, but for now, you must be in the Durable
26// Objects beta to be able to use the new syntax, while we work out the quirks.
27//
28// To see an example configuration for uploading module-based Workers, check out the wrangler.toml
29// file or one of our Durable Object templates for Wrangler:
30// * https://github.com/cloudflare/durable-objects-template
31// * https://github.com/cloudflare/durable-objects-rollup-esm
32// * https://github.com/cloudflare/durable-objects-webpack-commonjs
33 
34// ===============================
35// Required Environment
36// ===============================
37//
38// This worker, when deployed, must be configured with two environment bindings:
39// * rooms: A Durable Object namespace binding mapped to the ChatRoom class.
40// * limiters: A Durable Object namespace binding mapped to the RateLimiter class.
41//
42// Incidentally, in pre-modules Workers syntax, "bindings" (like KV bindings, secrets, etc.)
43// appeared in your script as global variables, but in the new modules syntax, this is no longer
44// the case. Instead, bindings are now delivered in an "environment object" when an event handler
45// (or Durable Object class constructor) is called. Look for the variable `env` below.
46//
47// We made this change, again, for composability: The global scope is global, but if you want to
48// call into existing code that has different environment requirements, then you need to be able
49// to pass the environment as a parameter instead.
50//
51// Once again, see the wrangler.toml file to understand how the environment is configured.
52 
53// =======================================================================================
54// The regular Worker part...
55//
56// This section of the code implements a normal Worker that receives HTTP requests from external
57// clients. This part is stateless.
58 
59// With the introduction of modules, we're experimenting with allowing text/data blobs to be
60// uploaded and exposed as synthetic modules. In wrangler.toml we specify a rule that files ending
61// in .html should be uploaded as "Data", equivalent to content-type `application/octet-stream`.
62// So when we import it as `HTML` here, we get the HTML content as an `ArrayBuffer`. This lets us
63// serve our app's static asset without relying on any separate storage. (However, the space
64// available for assets served this way is very limited; larger sites should continue to use Workers
65// KV to serve assets.)
66import HTML from "./chat.html";
67 
68// `handleErrors()` is a little utility function that can wrap an HTTP request handler in a
69// try/catch and return errors to the client. You probably wouldn't want to use this in production
70// code but it is convenient when debugging and iterating.
71async function handleErrors(request, func) {
72 try {
73 return await func();
74 } catch (err) {
75 if (request.headers.get("Upgrade") == "websocket") {
76 // Annoyingly, if we return an HTTP error in response to a WebSocket request, Chrome devtools
77 // won't show us the response body! So... let's send a WebSocket response with an error
78 // frame instead.
79 let pair = new WebSocketPair();
80 pair[1].accept();
81 pair[1].send(JSON.stringify({error: err.stack}));
82 pair[1].close(1011, "Uncaught exception during session setup");
83 return new Response(null, { status: 101, webSocket: pair[0] });
84 } else {
85 return new Response(err.stack, {status: 500});
86 }
87 }
88}
89 
90// In modules-syntax workers, we use `export default` to export our script's main event handlers.
91// Here, we export one handler, `fetch`, for receiving HTTP requests. In pre-modules workers, the
92// fetch handler was registered using `addEventHandler("fetch", event => { ... })`; this is just
93// new syntax for essentially the same thing.
94//
95// `fetch` isn't the only handler. If your worker runs on a Cron schedule, it will receive calls
96// to a handler named `scheduled`, which should be exported here in a similar way. We will be
97// adding other handlers for other types of events over time.
98export default {
99 async fetch(request, env) {
100 return await handleErrors(request, async () => {
101 // We have received an HTTP request! Parse the URL and route the request.
102 
103 let url = new URL(request.url);
104 let path = url.pathname.slice(1).split('/');
105 
106 if (!path[0]) {
107 // Serve our HTML at the root path.
108 return new Response(HTML, {headers: {"Content-Type": "text/html;charset=UTF-8"}});
109 }
110 
111 switch (path[0]) {
112 case "api":
113 // This is a request for `/api/...`, call the API handler.
114 return handleApiRequest(path.slice(1), request, env);
115 
116 default:
117 return new Response("Not found", {status: 404});
118 }
119 });
120 }
121}
122 
123 
124async function handleApiRequest(path, request, env) {
125 // We've received at API request. Route the request based on the path.
126 
127 switch (path[0]) {
128 case "room": {
129 // Request for `/api/room/...`.
130 
131 if (!path[1]) {
132 // The request is for just "/api/room", with no ID.
133 if (request.method == "POST") {
134 // POST to /api/room creates a private room.
135 //
136 // Incidentally, this code doesn't actually store anything. It just generates a valid
137 // unique ID for this namespace. Each durable object namespace has its own ID space, but
138 // IDs from one namespace are not valid for any other.
139 //
140 // The IDs returned by `newUniqueId()` are unguessable, so are a valid way to implement
141 // "anyone with the link can access" sharing. Additionally, IDs generated this way have
142 // a performance benefit over IDs generated from names: When a unique ID is generated,
143 // the system knows it is unique without having to communicate with the rest of the
144 // world -- i.e., there is no way that someone in the UK and someone in New Zealand
145 // could coincidentally create the same ID at the same time, because unique IDs are,
146 // well, unique!
147 let id = env.rooms.newUniqueId();
148 return new Response(id.toString(), {headers: {"Access-Control-Allow-Origin": "*"}});
149 } else {
150 // If we wanted to support returning a list of public rooms, this might be a place to do
151 // it. The list of room names might be a good thing to store in KV, though a singleton
152 // Durable Object is also a possibility as long as the Cache API is used to cache reads.
153 // (A caching layer would be needed because a single Durable Object is single-threaded,
154 // so the amount of traffic it can handle is limited. Also, caching would improve latency
155 // for users who don't happen to be located close to the singleton.)
156 //
157 // For this demo, though, we're not implementing a public room list, mainly because
158 // inevitably some trolls would probably register a bunch of offensive room names. Sigh.
159 return new Response("Method not allowed", {status: 405});
160 }
161 }
162 
163 // OK, the request is for `/api/room/<name>/...`. It's time to route to the Durable Object
164 // for the specific room.
165 let name = path[1];
166 
167 // Each Durable Object has a 256-bit unique ID. IDs can be derived from string names, or
168 // chosen randomly by the system.
169 let id;
170 if (name.match(/^[0-9a-f]{64}$/)) {
171 // The name is 64 hex digits, so let's assume it actually just encodes an ID. We use this
172 // for private rooms. `idFromString()` simply parses the text as a hex encoding of the raw
173 // ID (and verifies that this is a valid ID for this namespace).
174 id = env.rooms.idFromString(name);
175 } else if (name.length <= 32) {
176 // Treat as a string room name (limited to 32 characters). `idFromName()` consistently
177 // derives an ID from a string.
178 id = env.rooms.idFromName(name);
179 } else {
180 return new Response("Name too long", {status: 404});
181 }
182 
183 // Get the Durable Object stub for this room! The stub is a client object that can be used
184 // to send messages to the remote Durable Object instance. The stub is returned immediately;
185 // there is no need to await it. This is important because you would not want to wait for
186 // a network round trip before you could start sending requests. Since Durable Objects are
187 // created on-demand when the ID is first used, there's nothing to wait for anyway; we know
188 // an object will be available somewhere to receive our requests.
189 let roomObject = env.rooms.get(id);
190 
191 // Compute a new URL with `/api/room/<name>` removed. We'll forward the rest of the path
192 // to the Durable Object.
193 let newUrl = new URL(request.url);
194 newUrl.pathname = "/" + path.slice(2).join("/");
195 
196 // Send the request to the object. The `fetch()` method of a Durable Object stub has the
197 // same signature as the global `fetch()` function, but the request is always sent to the
198 // object, regardless of the request's URL.
199 return roomObject.fetch(newUrl, request);
200 }
201 
202 default:
203 return new Response("Not found", {status: 404});
204 }
205}
206 
207// =======================================================================================
208// The ChatRoom Durable Object Class
209 
210// ChatRoom implements a Durable Object that coordinates an individual chat room. Participants
211// connect to the room using WebSockets, and the room broadcasts messages from each participant
212// to all others.
213export class ChatRoom {
214 constructor(controller, env) {
215 // `controller.storage` provides access to our durable storage. It provides a simple KV
216 // get()/put() interface.
217 this.storage = controller.storage;
218 
219 // `env` is our environment bindings (discussed earlier).
220 this.env = env;
221 
222 // We will put the WebSocket objects for each client, along with some metadata, into
223 // `sessions`.
224 this.sessions = [];
225 
226 // We keep track of the last-seen message's timestamp just so that we can assign monotonically
227 // increasing timestamps even if multiple messages arrive simultaneously (see below). There's
228 // no need to store this to disk since we assume if the object is destroyed and recreated, much
229 // more than a millisecond will have gone by.
230 this.lastTimestamp = 0;
231 }
232 
233 // The system will call fetch() whenever an HTTP request is sent to this Object. Such requests
234 // can only be sent from other Worker code, such as the code above; these requests don't come
235 // directly from the internet. In the future, we will support other formats than HTTP for these
236 // communications, but we started with HTTP for its familiarity.
237 async fetch(request) {
238 return await handleErrors(request, async () => {
239 let url = new URL(request.url);
240 
241 switch (url.pathname) {
242 case "/websocket": {
243 // The request is to `/api/room/<name>/websocket`. A client is trying to establish a new
244 // WebSocket session.
245 if (request.headers.get("Upgrade") != "websocket") {
246 return new Response("expected websocket", {status: 400});
247 }
248 
249 // Get the client's IP address for use with the rate limiter.
250 let ip = request.headers.get("CF-Connecting-IP");
251 
252 // To accept the WebSocket request, we create a WebSocketPair (which is like a socketpair,
253 // i.e. two WebSockets that talk to each other), we return one end of the pair in the
254 // response, and we operate on the other end. Note that this API is not part of the
255 // Fetch API standard; unfortunately, the Fetch API / Service Workers specs do not define
256 // any way to act as a WebSocket server today.
257 let pair = new WebSocketPair();
258 
259 // We're going to take pair[1] as our end, and return pair[0] to the client.
260 await this.handleSession(pair[1], ip);
261 
262 // Now we return the other end of the pair to the client.
263 return new Response(null, { status: 101, webSocket: pair[0] });
264 }
265 
266 default:
267 return new Response("Not found", {status: 404});
268 }
269 });
270 }
271 
272 // handleSession() implements our WebSocket-based chat protocol.
273 async handleSession(webSocket, ip) {
274 // Accept our end of the WebSocket. This tells the runtime that we'll be terminating the
275 // WebSocket in JavaScript, not sending it elsewhere.
276 webSocket.accept();
277 
278 // Set up our rate limiter client.
279 let limiterId = this.env.limiters.idFromName(ip);
280 let limiter = new RateLimiterClient(
281 () => this.env.limiters.get(limiterId),
282 err => webSocket.close(1011, err.stack));
283 
284 // Create our session and add it to the sessions list.
285 // We don't send any messages to the client until it has sent us the initial user info
286 // message. Until then, we will queue messages in `session.blockedMessages`.
287 let session = {webSocket, blockedMessages: []};
288 this.sessions.push(session);
289 
290 // Queue "join" messages for all online users, to populate the client's roster.
291 this.sessions.forEach(otherSession => {
292 if (otherSession.name) {
293 session.blockedMessages.push(JSON.stringify({joined: otherSession.name}));
294 }
295 });
296 
297 // Load the last 100 messages from the chat history stored on disk, and send them to the
298 // client.
299 let storage = await this.storage.list({reverse: true, limit: 100});
300 let backlog = [...storage.values()];
301 backlog.reverse();
302 backlog.forEach(value => {
303 session.blockedMessages.push(value);
304 });
305 
306 // Set event handlers to receive messages.
307 let receivedUserInfo = false;
308 webSocket.addEventListener("message", async msg => {
309 try {
310 if (session.quit) {
311 // Whoops, when trying to send to this WebSocket in the past, it threw an exception and
312 // we marked it broken. But somehow we got another message? I guess try sending a
313 // close(), which might throw, in which case we'll try to send an error, which will also
314 // throw, and whatever, at least we won't accept the message. (This probably can't
315 // actually happen. This is defensive coding.)
316 webSocket.close(1011, "WebSocket broken.");
317 return;
318 }
319 
320 // Check if the user is over their rate limit and reject the message if so.
321 if (!limiter.checkLimit()) {
322 webSocket.send(JSON.stringify({
323 error: "Your IP is being rate-limited, please try again later."
324 }));
325 return;
326 }
327 
328 // I guess we'll use JSON.
329 let data = JSON.parse(msg.data);
330 
331 if (!receivedUserInfo) {
332 // The first message the client sends is the user info message with their name. Save it
333 // into their session object.
334 session.name = "" + (data.name || "anonymous");
335 
336 // Don't let people use ridiculously long names. (This is also enforced on the client,
337 // so if they get here they are not using the intended client.)
338 if (session.name.length > 32) {
339 webSocket.send(JSON.stringify({error: "Name too long."}));
340 webSocket.close(1009, "Name too long.");
341 return;
342 }
343 
344 // Deliver all the messages we queued up since the user connected.
345 session.blockedMessages.forEach(queued => {
346 webSocket.send(queued);
347 });
348 delete session.blockedMessages;
349 
350 // Broadcast to all other connections that this user has joined.
351 this.broadcast({joined: session.name});
352 
353 webSocket.send(JSON.stringify({ready: true}));
354 
355 // Note that we've now received the user info message.
356 receivedUserInfo = true;
357 
358 return;
359 }
360 
361 // Construct sanitized message for storage and broadcast.
362 data = { name: session.name, message: "" + data.message };
363 
364 // Block people from sending overly long messages. This is also enforced on the client,
365 // so to trigger this the user must be bypassing the client code.
366 if (data.message.length > 256) {
367 webSocket.send(JSON.stringify({error: "Message too long."}));
368 return;
369 }
370 
371 // Add timestamp. Here's where this.lastTimestamp comes in -- if we receive a bunch of
372 // messages at the same time (or if the clock somehow goes backwards????), we'll assign
373 // them sequential timestamps, so at least the ordering is maintained.
374 data.timestamp = Math.max(Date.now(), this.lastTimestamp + 1);
375 this.lastTimestamp = data.timestamp;
376 
377 // Broadcast the message to all other WebSockets.
378 let dataStr = JSON.stringify(data);
379 this.broadcast(dataStr);
380 
381 // Save message.
382 let key = new Date(data.timestamp).toISOString();
383 await this.storage.put(key, dataStr);
384 } catch (err) {
385 // Report any exceptions directly back to the client. As with our handleErrors() this
386 // probably isn't what you'd want to do in production, but it's convenient when testing.
387 webSocket.send(JSON.stringify({error: err.stack}));
388 }
389 });
390 
391 // On "close" and "error" events, remove the WebSocket from the sessions list and broadcast
392 // a quit message.
393 let closeOrErrorHandler = evt => {
394 session.quit = true;
395 this.sessions = this.sessions.filter(member => member !== session);
396 if (session.name) {
397 this.broadcast({quit: session.name});
398 }
399 };
400 webSocket.addEventListener("close", closeOrErrorHandler);
401 webSocket.addEventListener("error", closeOrErrorHandler);
402 }
403 
404 // broadcast() broadcasts a message to all clients.
405 broadcast(message) {
406 // Apply JSON if we weren't given a string to start with.
407 if (typeof message !== "string") {
408 message = JSON.stringify(message);
409 }
410 
411 // Iterate over all the sessions sending them messages.
412 let quitters = [];
413 this.sessions = this.sessions.filter(session => {
414 if (session.name) {
415 try {
416 session.webSocket.send(message);
417 return true;
418 } catch (err) {
419 // Whoops, this connection is dead. Remove it from the list and arrange to notify
420 // everyone below.
421 session.quit = true;
422 quitters.push(session);
423 return false;
424 }
425 } else {
426 // This session hasn't sent the initial user info message yet, so we're not sending them
427 // messages yet (no secret lurking!). Queue the message to be sent later.
428 session.blockedMessages.push(message);
429 return true;
430 }
431 });
432 
433 quitters.forEach(quitter => {
434 if (quitter.name) {
435 this.broadcast({quit: quitter.name});
436 }
437 });
438 }
439}
440 
441// =======================================================================================
442// The RateLimiter Durable Object class.
443 
444// RateLimiter implements a Durable Object that tracks the frequency of messages from a particular
445// source and decides when messages should be dropped because the source is sending too many
446// messages.
447//
448// We utilize this in ChatRoom, above, to apply a per-IP-address rate limit. These limits are
449// global, i.e. they apply across all chat rooms, so if a user spams one chat room, they will find
450// themselves rate limited in all other chat rooms simultaneously.
451export class RateLimiter {
452 constructor(controller, env) {
453 // Timestamp at which this IP will next be allowed to send a message. Start in the distant
454 // past, i.e. the IP can send a message now.
455 this.nextAllowedTime = 0;
456 }
457 
458 // Our protocol is: POST when the IP performs an action, or GET to simply read the current limit.
459 // Either way, the result is the number of seconds to wait before allowing the IP to perform its
460 // next action.
461 async fetch(request) {
462 return await handleErrors(request, async () => {
463 let now = Date.now() / 1000;
464 
465 this.nextAllowedTime = Math.max(now, this.nextAllowedTime);
466 
467 if (request.method == "POST") {
468 // POST request means the user performed an action.
469 // We allow one action per 5 seconds.
470 this.nextAllowedTime += 5;
471 }
472 
473 // Return the number of seconds that the client needs to wait.
474 //
475 // We provide a "grace" period of 20 seconds, meaning that the client can make 4-5 requests
476 // in a quick burst before they start being limited.
477 let cooldown = Math.max(0, this.nextAllowedTime - now - 20);
478 return new Response(cooldown);
479 })
480 }
481}
482 
483// RateLimiterClient implements rate limiting logic on the caller's side.
484class RateLimiterClient {
485 // The constructor takes two functions:
486 // * getLimiterStub() returns a new Durable Object stub for the RateLimiter object that manages
487 // the limit. This may be called multiple times as needed to reconnect, if the connection is
488 // lost.
489 // * reportError(err) is called when something goes wrong and the rate limiter is broken. It
490 // should probably disconnect the client, so that they can reconnect and start over.
491 constructor(getLimiterStub, reportError) {
492 this.getLimiterStub = getLimiterStub;
493 this.reportError = reportError;
494 
495 // Call the callback to get the initial stub.
496 this.limiter = getLimiterStub();
497 
498 // When `inCooldown` is true, the rate limit is currently applied and checkLimit() will return
499 // false.
500 this.inCooldown = false;
501 }
502 
503 // Call checkLimit() when a message is received to decide if it should be blocked due to the
504 // rate limit. Returns `true` if the message should be accepted, `false` to reject.
505 checkLimit() {
506 if (this.inCooldown) {
507 return false;
508 }
509 this.inCooldown = true;
510 this.callLimiter();
511 return true;
512 }
513 
514 // callLimiter() is an internal method which talks to the rate limiter.
515 async callLimiter() {
516 try {
517 let response;
518 try {
519 // Currently, fetch() needs a valid URL even though it's not actually going to the
520 // internet. We may loosen this in the future to accept an arbitrary string. But for now,
521 // we have to provide a dummy URL that will be ignored at the other end anyway.
522 response = await this.limiter.fetch("https://dummy-url", {method: "POST"});
523 } catch (err) {
524 // `fetch()` threw an exception. This is probably because the limiter has been
525 // disconnected. Stubs implement E-order semantics, meaning that calls to the same stub
526 // are delivered to the remote object in order, until the stub becomes disconnected, after
527 // which point all further calls fail. This guarantee makes a lot of complex interaction
528 // patterns easier, but it means we must be prepared for the occasional disconnect, as
529 // networks are inherently unreliable.
530 //
531 // Anyway, get a new limiter and try again. If it fails again, something else is probably
532 // wrong.
533 this.limiter = this.getLimiterStub();
534 response = await this.limiter.fetch("https://dummy-url", {method: "POST"});
535 }
536 
537 // The response indicates how long we want to pause before accepting more requests.
538 let cooldown = +(await response.text());
539 await new Promise(resolve => setTimeout(resolve, cooldown * 1000));
540 
541 // Done waiting.
542 this.inCooldown = false;
543 } catch (err) {
544 this.reportError(err);
545 }
546 }
547}