hub.tsannotatedhub.tssource94 lines · 3.5 KB · raw
1// The MCP hub: one Durable Object holding every open REPL tab's WebSocket,
2// keyed by the tab's session id. A tool call becomes a command sent to one
3// tab, and the tab's reply resolves it.
4//
5// State is in memory. When `wrangler dev` reloads the Worker (any save
6// under worker/src), the sockets drop and tabs reconnect within seconds
7// (repl/useWebSocketMCP.jsx), so a tool change needs no page reload.
8import { DurableObject } from 'cloudflare:workers';
9import { replyMs, type SharedCommand } from './tab-hub-core';
10
11// The page's own commands, and the shared tools' (tab-hub-core.ts's, which
12// the dev tab answers as the hosted one does: website/src/jev/tabTools.mjs).
13export type Command =
14  | { type: 'play'; code: string }
15  | { type: 'stop' }
16  | { type: 'get-code' }
17  | { type: 'get-logs' }
18  | SharedCommand;
19
20// As the tab sent it: the dev tools read these fields, and the shared
21// tools rebuild the rest first (tab-hub-core.ts's rebuildReply).
22export type Reply = {
23  ok?: boolean;
24  error?: string;
25  code?: string;
26  logs?: { at: number; message: string; count: number }[];
27  [field: string]: unknown;
28};
29const SESSION = /^[a-z0-9-]{1,32}$/;
30
31export class Hub extends DurableObject {
32  private tabs = new Map<string, WebSocket>();
33  private pending = new Map<string, (reply: Reply) => void>();
34  private nextId = 0;
35
36  async fetch(request: Request): Promise<Response> {
37    const session = new URL(request.url).searchParams.get('session_id') ?? '';
38    if (request.headers.get('Upgrade') !== 'websocket') return new Response('WebSocket only', { status: 426 });
39    if (!SESSION.test(session)) return new Response('session_id: 1-32 of a-z, 0-9, -', { status: 400 });
40
41    const [client, server] = Object.values(new WebSocketPair());
42    server.accept();
43    this.tabs.get(session)?.close(1000, 'replaced by a newer connection');
44    this.tabs.set(session, server);
45    server.addEventListener('message', (event) => {
46      let message: { type?: string; id?: string } & Reply;
47      try {
48        message = JSON.parse(String(event.data));
49      } catch {
50        return;
51      }
52      if (message.type === 'reply' && message.id) {
53        this.pending.get(message.id)?.(message);
54        this.pending.delete(message.id);
55      }
56    });
57    const drop = () => {
58      if (this.tabs.get(session) === server) this.tabs.delete(session);
59    };
60    server.addEventListener('close', drop);
61    server.addEventListener('error', drop);
62    return new Response(null, { status: 101, webSocket: client });
63  }
64
65  sessions(): string[] {
66    return [...this.tabs.keys()];
67  }
68
69  async command(session: string, command: Command): Promise<Reply> {
70    const tab = this.tabs.get(session);
71    if (!tab) throw new Error(`no tab with session ${session}`);
72    const id = `${Date.now()}-${this.nextId++}`;
73    const ms = replyMs(command);
74    return new Promise<Reply>((resolve, reject) => {
75      const timer = setTimeout(() => {
76        this.pending.delete(id);
77        reject(new Error(`tab ${session} did not answer in ${ms / 1000} s`));
78      }, ms);
79      this.pending.set(id, (reply) => {
80        clearTimeout(timer);
81        resolve(reply);
82      });
83      try {
84        tab.send(JSON.stringify({ type: 'command', id, command }));
85      } catch {
86        // the tab went away between the lookup and the send (a reload)
87        clearTimeout(timer);
88        this.pending.delete(id);
89        if (this.tabs.get(session) === tab) this.tabs.delete(session);
90        reject(new Error(`tab ${session} disconnected`));
91      }
92    });
93  }
94}