jevstrudel.git / worker / src / tab-hub.ts

One signed-in user's tab hub: a Durable Object per account (TAB_HUB.idFromName(userId)) holding the WebSockets of that user's open REPL tabs, so the hosted MCP (hosted-mcp.ts), acting for that user with an OAuth token, can play code in them. A tab joins only through GET /jev/me/tabs (index.ts), which takes the session cookie from this site's own pages and nowhere else, so the only tabs in a user's hub are that user's; and the MCP reaches only the hub its token's user names. No other visitor's tab is ever in reach.

Unlike the dev hub (hub.ts), it runs in production, so it uses the Hibernation API: a tab sitting open between commands costs no duration. A command in flight keeps the object awake (the RPC call is waiting), so the replies it waits for are held in memory. The rules (sizes, which tab, 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';
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  }

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

The user's open tabs, by session id.

94  sessions(): string[] {
95    return [...new Set(this.tabs().map((t) => t.at.session))];
96  }
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}