jevstrudel.git / worker / src / tab-hub.ts
1// One signed-in user's tab hub: a Durable Object per account
2// (`TAB_HUB.idFromName(userId)`) holding the WebSockets of that user's open
3// REPL tabs, so the hosted MCP (hosted-mcp.ts), acting for that user with
4// an OAuth token, can play code in them. A tab joins only through
5// `GET /jev/me/tabs` (index.ts), which takes the session cookie from this
6// site's own pages and nowhere else, so the only tabs in a user's hub are
7// that user's; and the MCP reaches only the hub its token's user names. No
8// other visitor's tab is ever in reach.
9//
10// Unlike the dev hub (hub.ts), it runs in production, so it uses the
11// Hibernation API: a tab sitting open between commands costs no duration.
12// A command in flight keeps the object awake (the RPC call is waiting), so
13// the replies it waits for are held in memory. The rules (sizes, which tab,
14// what a reply may hold) are tab-hub-core.ts.
15import { DurableObject } from 'cloudflare:workers';
16import { MAX_FRAME_BYTES, MAX_TABS, parseReply, replyMs, SESSION, type TabCommand, type TabReply } from './tab-hub-core';
17
18type Attachment = { session: string; joined: number };
19type Waiting = { socket: WebSocket; resolve: (reply: TabReply) => void };
20
21export class TabHub extends DurableObject<object> {
22  private waiting = new Map<string, Waiting>();
23  private nextId = 0;
24
25  constructor(ctx: DurableObjectState, env: object) {
26    super(ctx, env);
27    // keeps a tab's connection alive without waking the object
28    ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair('ping', 'pong'));
29  }
30
31  private tabs(): { socket: WebSocket; at: Attachment }[] {
32    return this.ctx
33      .getWebSockets()
34      .map((socket) => ({ socket, at: socket.deserializeAttachment() as Attachment | null }))
35      .filter((t): t is { socket: WebSocket; at: Attachment } => !!t.at);
36  }
37
38  // A tab joining: index.ts has checked who is signed in, and routed here by their id.
39  async fetch(request: Request): Promise<Response> {
40    if (request.headers.get('Upgrade') !== 'websocket') return new Response('WebSocket only', { status: 426 });
41    const session = new URL(request.url).searchParams.get('session_id') ?? '';
42    if (!SESSION.test(session)) return new Response('session_id: 1-32 of a-z, 0-9, -', { status: 400 });
43
44    const tabs = this.tabs();
45    // the same tab reconnecting replaces itself; past MAX_TABS the oldest goes
46    for (const t of tabs.filter((t) => t.at.session === session)) t.socket.close(1000, 'replaced by a newer connection');
47    const others = tabs.filter((t) => t.at.session !== session).sort((a, b) => a.at.joined - b.at.joined);
48    for (const t of others.slice(0, Math.max(0, others.length - (MAX_TABS - 1)))) {
49      t.socket.close(4008, 'too many tabs; the oldest was closed');
50    }
51
52    const [client, server] = Object.values(new WebSocketPair());
53    this.ctx.acceptWebSocket(server, [session]);
54    server.serializeAttachment({ session, joined: Date.now() } satisfies Attachment);
55    return new Response(null, { status: 101, webSocket: client });
56  }
57
58  webSocketMessage(socket: WebSocket, message: string | ArrayBuffer): void {
59    const size = typeof message === 'string' ? message.length : message.byteLength;
60    if (size > MAX_FRAME_BYTES) {
61      socket.close(1009, 'message too large');
62      return;
63    }
64    const parsed = parseReply(message);
65    if (!parsed) return;
66    const w = this.waiting.get(parsed.id);
67    // only the tab the command went to may answer it
68    if (!w || w.socket !== socket) return;
69    this.waiting.delete(parsed.id);
70    w.resolve(parsed.reply);
71  }
72
73  webSocketClose(socket: WebSocket, code: number): void {
74    this.dropWaiting(socket);
75    try {
76      socket.close(code === 1005 || code === 1006 ? 1000 : code, 'closed');
77    } catch {}
78  }
79
80  webSocketError(socket: WebSocket): void {
81    this.dropWaiting(socket);
82  }
83
84  private dropWaiting(socket: WebSocket) {
85    for (const [id, w] of this.waiting) {
86      if (w.socket === socket) {
87        this.waiting.delete(id);
88        w.resolve({ error: 'the tab closed before it answered' });
89      }
90    }
91  }
92
93  // The user's open tabs, by session id.
94  sessions(): string[] {
95    return [...new Set(this.tabs().map((t) => t.at.session))];
96  }
97
98  async command(session: string, command: TabCommand): Promise<TabReply> {
99    const tab = this.tabs().find((t) => t.at.session === session);
100    if (!tab) throw new Error(`you have no open tab ${session}`);
101    const id = `${Date.now()}-${this.nextId++}`;
102    const ms = replyMs(command);
103    return new Promise<TabReply>((resolve, reject) => {
104      const timer = setTimeout(() => {
105        this.waiting.delete(id);
106        reject(new Error(`tab ${session} did not answer in ${ms / 1000} s`));
107      }, ms);
108      this.waiting.set(id, {
109        socket: tab.socket,
110        resolve: (reply) => {
111          clearTimeout(timer);
112          resolve(reply);
113        },
114      });
115      try {
116        tab.socket.send(JSON.stringify({ type: 'command', id, command }));
117      } catch {
118        clearTimeout(timer);
119        this.waiting.delete(id);
120        reject(new Error(`tab ${session} disconnected`));
121      }
122    });
123  }
124}