Skip to content
115 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Chat and workspace agents: channels, DMs and named agents you talk to1/**
2 * One channel's live room, as a Durable Object named by the channel id.
3 *
4 * Why an object per channel: everyone looking at a channel must see the
5 * same messages in the same order the moment they are written, and one
6 * object per channel sees every socket open on it. It holds the sockets
7 * with the WebSocket Hibernation API, so a quiet channel with people
8 * looking at it costs nothing between messages: the object is evicted
9 * from memory and the sockets stay open at the edge. Each socket carries
10 * who is behind it (`serializeAttachment`), which survives hibernation.
11 *
12 * The room authorizes nothing. The Worker checks a viewer may read the
13 * channel before forwarding their socket here (src/index.ts, `live`), and
14 * tells the room to drop someone's sockets when they leave a private one.
15 */
16import { DurableObject } from "cloudflare:workers";
17
18import { principalKey, type ChatLiveEvent, type MemberProfile } from "@g1t/contracts";
19
20import { PERSON_TYPING_MS, isTypingFrame } from "./messages.ts";
21
22/** Header the Worker sets on a socket it forwards: who it is, as JSON (`RoomMember`). */
23export const ROOM_MEMBER_HEADER = "x-g1t-chat-member";
24
25export type RoomMember = { channel_id: string; member: MemberProfile };
26
27/** The shortest time between two "is typing" events from one socket. */
28const TYPING_EVERY_MS = 2_000;
29
30export class ChannelRoom extends DurableObject<object> {
31 /** When each socket last said it was typing; lost on hibernation, which only lets one more through. */
32 private typedAt = new WeakMap<WebSocket, number>();
33
34 constructor(ctx: DurableObjectState, env: object) {
35 super(ctx, env);
36 // Keepalives are answered without waking the object.
37 ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair("ping", "pong"));
38 }
39
40 override async fetch(request: Request): Promise<Response> {
41 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
42 return new Response("Expected a WebSocket upgrade\n", { status: 426 });
43 }
44 let who: RoomMember;
45 try {
46 who = JSON.parse(request.headers.get(ROOM_MEMBER_HEADER) ?? "") as RoomMember;
47 } catch {
48 return new Response("Missing member\n", { status: 400 });
49 }
50 const pair = new WebSocketPair();
51 const [client, server] = Object.values(pair) as [WebSocket, WebSocket];
52 // Tagged with the member's key, so their sockets can be found to drop.
53 this.ctx.acceptWebSocket(server, [principalKey(who.member)]);
54 server.serializeAttachment(who);
55 return new Response(null, { status: 101, webSocket: client });
56 }
57
58 /** Sends `event` to every socket in the room, except those of `except` (a principal key). */
59 broadcast(event: ChatLiveEvent, except: string | null = null): void {
60 const text = JSON.stringify(event);
61 for (const socket of this.ctx.getWebSockets()) {
62 if (except && this.ctx.getTags(socket).includes(except)) continue;
63 try {
64 socket.send(text);
65 } catch {
66 // Closing already; webSocketClose tidies up.
67 }
68 }
69 }
70
71 /** Closes a member's sockets: they left a private channel and may no longer read it. */
72 drop(principal: string): void {
73 for (const socket of this.ctx.getWebSockets(principal)) {
74 try {
75 socket.close(4403, "No longer a member");
76 } catch {
77 // Already closed.
78 }
79 }
80 }
81
82 override async webSocketMessage(socket: WebSocket, message: string | ArrayBuffer): Promise<void> {
83 const who = socket.deserializeAttachment() as RoomMember | null;
84 if (!who) return;
85 // The only thing a client says over the socket: it is typing
86 // (`{"type":"typing","channel_id":"<id>"}`, at most every 3 s).
87 // Everything else (posting, reading) goes through the site, which
88 // checks it.
89 if (!isTypingFrame(message, who.channel_id)) return;
90 const now = Date.now();
91 if (now - (this.typedAt.get(socket) ?? 0) < TYPING_EVERY_MS) return;
92 this.typedAt.set(socket, now);
93 this.broadcast(
94 {
95 type: "typing",
96 channel_id: who.channel_id,
97 member: who.member,
98 until: new Date(now + PERSON_TYPING_MS).toISOString(),
99 },
100 principalKey(who.member),
101 );
102 }
103
104 override async webSocketClose(socket: WebSocket, code: number, reason: string): Promise<void> {
105 try {
106 socket.close(code, reason);
107 } catch {
108 // Already closed.
109 }
110 }
111
112 override async webSocketError(): Promise<void> {
113 // Nothing to keep: the socket is gone, and getWebSockets no longer lists it.
114 }
115}