diff --git a/.changeset/calm-sockets-share.md b/.changeset/calm-sockets-share.md new file mode 100644 index 0000000..61b71ce --- /dev/null +++ b/.changeset/calm-sockets-share.md @@ -0,0 +1,5 @@ +--- +"@mcp-b/do-runtime": patch +--- + +Export the shared in-memory hibernation mirror and browser WebSocket-upgrade adapter so embedders can reuse the reference host behavior instead of copying example shims. diff --git a/README.md b/README.md index 0fb9de3..d6ce1e4 100644 --- a/README.md +++ b/README.md @@ -241,6 +241,12 @@ registry is populated before the constructor, so SDKs can lazily rebuild their connection wrappers without another upgrade or connect hook. Closed sockets are removed before `webSocketClose` runs. +`HibernationMirror` is the package's in-memory reference implementation. Seed a +replacement mirror with the prior socket snapshot and auto-response pair, then +pass that mirror to `ports.hibernation` and its `snapshot()` to `webSockets`. +Browser hosts can install the remaining Request/Response upgrade accommodation +from `@mcp-b/do-runtime/browser`; the runtime itself supplies `WebSocketPair`. + `container.quiescence()` reports armed timers, pending `waitUntil` work, input lock state, and output-gate breakage without waiting. `drainWaitUntil()` is for shutdown and intentionally never settles while a live interval remains armed. diff --git a/conformance/browser/actor.worker.ts b/conformance/browser/actor.worker.ts index 6034cd4..5c06aac 100644 --- a/conformance/browser/actor.worker.ts +++ b/conformance/browser/actor.worker.ts @@ -49,6 +49,7 @@ import { actorScopeBindings, createActorContainer, + HibernationMirror, installActorScope, newRpcSession, type ActorContainer, @@ -77,12 +78,13 @@ import { type SqliteWasmHost, } from "../../backends/sqlite-wasm"; import { Probe } from "../fixtures/probe"; -import { HibernationMirror } from "../../examples/platform-shims/hibernation-mirror"; import { installWebSocketUpgradeGlobals, upgradeWebSocket, - webSocketUpgradeRequest, type UpgradeWebSocket, +} from "../../src/browser"; +import { + webSocketUpgradeRequest, } from "../websocket-upgrade"; import type { ActorBoot, ActorRpc, SupervisorRpc } from "./protocol"; import { installPool, timer, UNIQUE_KEY } from "./substrate"; diff --git a/conformance/node/host.ts b/conformance/node/host.ts index 726ada7..f1ac603 100644 --- a/conformance/node/host.ts +++ b/conformance/node/host.ts @@ -35,6 +35,7 @@ import { createNodeSqlProvider } from "../../backends/node-sqlite"; import { AlarmScheduler, createActorContainer, + HibernationMirror, installWebSocketGlobals, type ActorContainer, type ActorEntry, @@ -64,12 +65,13 @@ import type { ProbeActor, } from "../host"; import { Probe } from "../fixtures/probe"; -import { HibernationMirror } from "../../examples/platform-shims/hibernation-mirror"; import { installWebSocketUpgradeGlobals, upgradeWebSocket, - webSocketUpgradeRequest, type UpgradeWebSocket, +} from "../../src/browser"; +import { + webSocketUpgradeRequest, } from "../websocket-upgrade"; // ======================================================================================= diff --git a/conformance/websocket-upgrade.ts b/conformance/websocket-upgrade.ts index bc72f02..9c1b99d 100644 --- a/conformance/websocket-upgrade.ts +++ b/conformance/websocket-upgrade.ts @@ -1,45 +1,7 @@ -import { markWebSocketUsed, type RawWebSocket } from "../src/index"; - -export type UpgradeWebSocket = RawWebSocket & { - accept(): void; - readonly readyState: number; -}; - -type UpgradeResponseInit = ResponseInit & { webSocket?: UpgradeWebSocket }; - -let installed = false; - -/** The Request/Response half of WebSocket upgrades supplied by workerd in the oracle lane. */ -export function installWebSocketUpgradeGlobals(): void { - if (installed) return; - installed = true; - const NativeResponse = globalThis.Response; - class WorkersResponse extends NativeResponse { - constructor(body?: BodyInit | null, init: UpgradeResponseInit = {}) { - const upgrade = init.status === 101; - const { webSocket, ...nativeInit } = init; - super(body, upgrade ? { ...nativeInit, status: 200 } : nativeInit); - if (upgrade) Object.defineProperty(this, "status", { value: 101 }); - if (webSocket !== undefined) { - markWebSocketUsed(webSocket); - Object.defineProperty(this, "webSocket", { value: webSocket }); - } - } - } - globalThis.Response = WorkersResponse as typeof Response; -} - -export function upgradeWebSocket(response: Response): UpgradeWebSocket | undefined { - return (response as Response & { webSocket?: UpgradeWebSocket }).webSocket; -} +import { withWebSocketUpgrade } from "../src/browser"; /** Preserve workerd's upgrade signal when browser Request.clone() drops forbidden headers. */ export function webSocketUpgradeRequest(url: string, tags: readonly string[]): Request { const request = new Request(`${url}?${tags.map((tag) => `tag=${encodeURIComponent(tag)}`).join("&")}`); - const get = request.headers.get.bind(request.headers); - Object.defineProperty(request.headers, "get", { - value: (name: string): string | null => - name.toLowerCase() === "upgrade" ? "websocket" : get(name), - }); - return request; + return withWebSocketUpgrade(request); } diff --git a/examples/extension/README.md b/examples/extension/README.md index 2792eed..8d91bdc 100644 --- a/examples/extension/README.md +++ b/examples/extension/README.md @@ -134,8 +134,8 @@ only as a competing supervisor and asserts that Web Locks refuse it before OPFS. | `src/background.ts` | The service worker: offscreen lifecycle and `chrome.alarms` projection. | | `src/popup/popup.ts` | Four buttons and an output pane. | | `src/protocol.ts` | The types both TypeScript projects compile. It imports nothing. | -| `../platform-shims/memory-websocket-pair.ts` | The browser `Response`-101 shim; the runtime supplies `WebSocketPair`. | -| `../platform-shims/hibernation-mirror.ts` | The process-local `HibernationHost` record used by this example and the conformance embedders. | +| `@mcp-b/do-runtime/browser` | The browser Request/`Response`-101 upgrade adapter; the runtime supplies `WebSocketPair`. | +| `@mcp-b/do-runtime` `HibernationMirror` | The process-local `HibernationHost` record shared by this example and the conformance embedders. | | `../platform-shims/message-port-websocket.ts` | The client-side WebSocket adapter carried over a `MessagePort`. | | `public/manifest.json` | Copied verbatim into `dist/` by Vite's `publicDir`. | diff --git a/examples/extension/src/worker/actor.worker.ts b/examples/extension/src/worker/actor.worker.ts index 9189c33..f3bbcd6 100644 --- a/examples/extension/src/worker/actor.worker.ts +++ b/examples/extension/src/worker/actor.worker.ts @@ -28,6 +28,7 @@ import { createActorContainer, createDurableObjectNamespace, gateRequestBody, + HibernationMirror, installActorScope, newRpcSession, type ActorContainer, @@ -51,12 +52,11 @@ import sqlite3InitModule from "@sqlite.org/sqlite-wasm"; import { getAgentByName, routeAgentEmail, routeAgentRequest } from "agents"; import { RpcTarget } from "cloudflare:workers"; import { - installWebSocketUpgradeResponse, + installWebSocketUpgradeGlobals, upgradeWebSocket, withWebSocketUpgrade, type UpgradeWebSocket, -} from "../../../platform-shims/memory-websocket-pair"; -import { HibernationMirror } from "../../../platform-shims/hibernation-mirror"; +} from "@mcp-b/do-runtime/browser"; import { serveMessagePortWebSockets } from "../../../platform-shims/message-port-websocket"; import type { CounterSnapshot, @@ -386,7 +386,7 @@ function counterNamespace(gate?: FacetModule["gate"]) { const rootNamespace = counterNamespace(); -installWebSocketUpgradeResponse(); +installWebSocketUpgradeGlobals(); /** * What the installed globals resolve to, and it REFUSES rather than falling diff --git a/examples/platform-shims/hibernation-mirror.ts b/examples/platform-shims/hibernation-mirror.ts deleted file mode 100644 index fa56f55..0000000 --- a/examples/platform-shims/hibernation-mirror.ts +++ /dev/null @@ -1,38 +0,0 @@ -import type { - HibernationHost, - RawWebSocket, - RehydratedWebSocket, -} from "@mcp-b/do-runtime"; - -type MirroredWebSocket = RehydratedWebSocket & { tags: readonly string[] }; - -/** In-memory socket state shared by the reference embedders and browser example. */ -export class HibernationMirror implements HibernationHost { - readonly #entries = new Map(); - autoResponsePair: { request: string; response: string } | null = null; - - accepted(socket: RawWebSocket, tags: readonly string[]): void { - this.#entries.set(socket, { socket, tags: [...tags] }); - } - - attachment(socket: RawWebSocket, bytes: Uint8Array | null): void { - const entry = this.#entries.get(socket); - if (entry === undefined) { - throw new Error("Hibernation mirror: attachment preceded socket acceptance."); - } - if (bytes === null) delete entry.attachment; - else entry.attachment = bytes.slice(); - } - - autoResponse(pair: { request: string; response: string } | null): void { - this.autoResponsePair = pair === null ? null : { ...pair }; - } - - closed(socket: RawWebSocket): void { - this.#entries.delete(socket); - } - - snapshot(): RehydratedWebSocket[] { - return [...this.#entries.values()]; - } -} diff --git a/examples/platform-shims/memory-websocket-pair.ts b/examples/platform-shims/memory-websocket-pair.ts deleted file mode 100644 index 21e4029..0000000 --- a/examples/platform-shims/memory-websocket-pair.ts +++ /dev/null @@ -1,44 +0,0 @@ -import { markWebSocketUsed, type RawWebSocket } from "@mcp-b/do-runtime"; - -export type UpgradeWebSocket = EventTarget & RawWebSocket & { - accept(): void; - readonly readyState: number; -}; - -type UpgradeResponseInit = ResponseInit & { webSocket?: UpgradeWebSocket }; - -/** Install the Response-101 half; `installActorScope` supplies the runtime's WebSocketPair. */ -export function installWebSocketUpgradeResponse(): void { - const NativeResponse = globalThis.Response; - class WorkersResponse extends NativeResponse { - constructor(body?: BodyInit | null, init: UpgradeResponseInit = {}) { - const upgrade = init.status === 101; - const { webSocket, ...nativeInit } = init; - super(body, upgrade ? { ...nativeInit, status: 200 } : nativeInit); - if (upgrade) Object.defineProperty(this, "status", { value: 101 }); - if (webSocket !== undefined) { - markWebSocketUsed(webSocket); - Object.defineProperty(this, "webSocket", { value: webSocket }); - } - } - } - globalThis.Response = WorkersResponse as typeof Response; -} - -export function upgradeWebSocket(response: Response): UpgradeWebSocket | undefined { - return (response as Response & { webSocket?: UpgradeWebSocket }).webSocket; -} - -/** Preserve workerd's WebSocket upgrade signal across browser `Request.clone()` calls. */ -export function withWebSocketUpgrade(request: T): T { - const get = request.headers.get.bind(request.headers); - Object.defineProperty(request.headers, "get", { - value: (name: string): string | null => - name.toLowerCase() === "upgrade" ? "websocket" : get(name), - }); - const clone = request.clone.bind(request); - Object.defineProperty(request, "clone", { - value: (): T => withWebSocketUpgrade(clone()), - }); - return request; -} diff --git a/package.json b/package.json index 4b2e4c3..0982e01 100644 --- a/package.json +++ b/package.json @@ -41,6 +41,10 @@ "types": "./dist/src/index.d.ts", "import": "./dist/index.js" }, + "./browser": { + "types": "./dist/src/browser.d.ts", + "import": "./dist/browser.js" + }, "./server/alarm-scheduler": { "types": "./dist/src/server/alarm-scheduler.d.ts", "import": "./dist/server/alarm-scheduler.js" diff --git a/scripts/build-package.mjs b/scripts/build-package.mjs index b3e4b53..b7a292f 100644 --- a/scripts/build-package.mjs +++ b/scripts/build-package.mjs @@ -16,6 +16,7 @@ await build({ lib: { entry: { index: new URL("src/index.ts", root).pathname, + browser: new URL("src/browser.ts", root).pathname, "server/alarm-scheduler": new URL("src/server/alarm-scheduler.ts", root).pathname, "backends/sqlite-wasm": new URL("backends/sqlite-wasm.ts", root).pathname, "backends/node-sqlite": new URL("backends/node-sqlite.ts", root).pathname, diff --git a/src/browser.test.ts b/src/browser.test.ts new file mode 100644 index 0000000..80e9869 --- /dev/null +++ b/src/browser.test.ts @@ -0,0 +1,40 @@ +import { afterAll, describe, expect, test } from "vitest"; +import type { RawWebSocket } from "./index"; +import { + installWebSocketUpgradeGlobals, + upgradeWebSocket, + withWebSocketUpgrade, +} from "./browser"; + +const NativeRequest = globalThis.Request; +const NativeResponse = globalThis.Response; + +afterAll(() => { + globalThis.Request = NativeRequest; + globalThis.Response = NativeResponse; +}); + +describe("browser WebSocket upgrade globals", () => { + test("carry status 101, its socket, and the upgrade marker through reconstructed requests", () => { + installWebSocketUpgradeGlobals(); + installWebSocketUpgradeGlobals(); + const request = withWebSocketUpgrade(new Request("https://example.test/socket")); + + expect(request.clone().headers.get("Upgrade")).toBe("websocket"); + expect(new Request(request).headers.get("Upgrade")).toBe("websocket"); + + const socket: RawWebSocket & EventTarget & { accept(): void; readonly readyState: number } = + Object.assign(new EventTarget(), { + accept() {}, + close() {}, + readyState: WebSocket.OPEN, + send() {}, + }); + // SAFETY: the browser adapter deliberately accepts the smaller RawWebSocket + // host seam; the ambient Workers type only spells this field as WebSocket. + const response = new Response(null, { status: 101, webSocket: socket as WebSocket }); + + expect(response.status).toBe(101); + expect(upgradeWebSocket(response)).toBe(socket); + }); +}); diff --git a/src/browser.ts b/src/browser.ts new file mode 100644 index 0000000..20a3e2d --- /dev/null +++ b/src/browser.ts @@ -0,0 +1,73 @@ +import { markWebSocketUsed, type RawWebSocket } from "./api/web-socket"; + +export type UpgradeWebSocket = EventTarget & + RawWebSocket & { + accept(): void; + readonly readyState: number; + }; + +type UpgradeResponseInit = Omit & { + webSocket?: UpgradeWebSocket; +}; +type UpgradeResponse = Response & { readonly webSocket?: UpgradeWebSocket }; +type CloneableRequest = { readonly headers: Headers; clone(): CloneableRequest }; + +const upgradeRequests = new WeakSet(); +let installed = false; + +/** Install the Request/Response half of browser-hosted WebSocket upgrades. */ +export function installWebSocketUpgradeGlobals(): void { + if (installed) return; + installed = true; + + const NativeRequest = globalThis.Request; + class WorkersRequest extends NativeRequest { + constructor(input: RequestInfo | URL, init?: RequestInit) { + const upgrade = input instanceof NativeRequest && isWebSocketUpgrade(input); + super(input, init); + if (upgrade) withWebSocketUpgrade(this); + } + } + globalThis.Request = WorkersRequest as typeof Request; + + const NativeResponse = globalThis.Response; + class WorkersResponse extends NativeResponse { + constructor(body?: BodyInit | null, init: UpgradeResponseInit = {}) { + const upgrade = init.status === 101; + const { webSocket, ...nativeInit } = init; + super(body, upgrade ? { ...nativeInit, status: 200 } : nativeInit); + if (upgrade) Object.defineProperty(this, "status", { value: 101 }); + if (webSocket !== undefined) { + markWebSocketUsed(webSocket); + Object.defineProperty(this, "webSocket", { value: webSocket }); + } + } + } + globalThis.Response = WorkersResponse as typeof Response; +} + +export function upgradeWebSocket(response: Response): UpgradeWebSocket | undefined { + return (response as UpgradeResponse).webSocket; +} + +/** Preserve the upgrade signal across browser `Request.clone()` calls. */ +export function withWebSocketUpgrade(request: T): T { + if (upgradeRequests.has(request)) return request; + upgradeRequests.add(request); + const get = request.headers.get.bind(request.headers); + Object.defineProperty(request.headers, "get", { + value: (name: string): string | null => + name.toLowerCase() === "upgrade" ? "websocket" : get(name), + }); + const clone = request.clone.bind(request); + Object.defineProperty(request, "clone", { + value: (): CloneableRequest => withWebSocketUpgrade(clone()), + }); + return request; +} + +function isWebSocketUpgrade(request: Request): boolean { + return ( + upgradeRequests.has(request) || request.headers.get("Upgrade")?.toLowerCase() === "websocket" + ); +} diff --git a/src/index.ts b/src/index.ts index f43390b..ff4fada 100644 --- a/src/index.ts +++ b/src/index.ts @@ -110,6 +110,8 @@ export { FACET_ALARM_UNIMPLEMENTED_MESSAGE, noFacets, } from "./server/actor-container"; +export type { HibernationAutoResponse } from "./server/hibernation-mirror"; +export { HibernationMirror } from "./server/hibernation-mirror"; export type { ActorChannelFactory, GlobalActorRequest } from "./api/actor"; export { createDurableObjectNamespace } from "./server/actor-namespace"; export { ACTOR_CLASS_SERIALIZATION_UNIMPLEMENTED_MESSAGE } from "./api/actor"; diff --git a/src/server/hibernation-mirror.test.ts b/src/server/hibernation-mirror.test.ts new file mode 100644 index 0000000..9159476 --- /dev/null +++ b/src/server/hibernation-mirror.test.ts @@ -0,0 +1,46 @@ +import { describe, expect, test } from "vitest"; +import { + HibernationMirror, + type RawWebSocket, + type RehydratedWebSocket, +} from "../index"; + +function rawWebSocket(): RawWebSocket { + return { + addEventListener() {}, + close() {}, + send() {}, + }; +} + +describe("HibernationMirror", () => { + test("hands copied socket metadata to a replacement host", () => { + const socket = rawWebSocket(); + const tags = ["connection-id", "room"]; + const attachment = new Uint8Array([1, 2, 3]); + const mirror = new HibernationMirror(); + + mirror.accepted(socket, tags); + mirror.attachment(socket, attachment); + mirror.autoResponse({ request: "ping", response: "pong" }); + tags[0] = "changed"; + attachment[0] = 9; + + const snapshot = mirror.snapshot(); + expect(snapshot).toEqual([ + { socket, tags: ["connection-id", "room"], attachment: new Uint8Array([1, 2, 3]) }, + ] satisfies RehydratedWebSocket[]); + expect(mirror.autoResponsePair).toEqual({ request: "ping", response: "pong" }); + + const replacement = new HibernationMirror(snapshot, mirror.autoResponsePair); + snapshot[0]!.tags = ["mutated"]; + snapshot[0]!.attachment![0] = 8; + expect(replacement.snapshot()).toEqual([ + { socket, tags: ["connection-id", "room"], attachment: new Uint8Array([1, 2, 3]) }, + ]); + expect(replacement.autoResponsePair).toEqual({ request: "ping", response: "pong" }); + + mirror.closed(socket); + expect(mirror.snapshot()).toEqual([]); + }); +}); diff --git a/src/server/hibernation-mirror.ts b/src/server/hibernation-mirror.ts new file mode 100644 index 0000000..6ead7fb --- /dev/null +++ b/src/server/hibernation-mirror.ts @@ -0,0 +1,68 @@ +import type { + HibernationHost, + RawWebSocket, + RehydratedWebSocket, +} from "../api/web-socket"; + +export type HibernationAutoResponse = { request: string; response: string }; + +type MirroredWebSocket = RehydratedWebSocket & { tags: readonly string[] }; + +/** In-memory socket state shared by embedders that replace live actor containers. */ +export class HibernationMirror implements HibernationHost { + readonly #entries = new Map(); + #autoResponsePair: HibernationAutoResponse | null; + + constructor( + rehydrated: readonly RehydratedWebSocket[] = [], + autoResponsePair: HibernationAutoResponse | null = null, + ) { + for (const value of rehydrated) this.#entries.set(value.socket, cloneEntry(value)); + this.#autoResponsePair = cloneAutoResponse(autoResponsePair); + } + + get autoResponsePair(): HibernationAutoResponse | null { + return cloneAutoResponse(this.#autoResponsePair); + } + + accepted(socket: RawWebSocket, tags: readonly string[]): void { + this.#entries.set(socket, { socket, tags: [...tags] }); + } + + attachment(socket: RawWebSocket, bytes: Uint8Array | null): void { + const entry = this.#entries.get(socket); + if (entry === undefined) { + throw new Error("Hibernation mirror: attachment preceded socket acceptance."); + } + if (bytes === null) delete entry.attachment; + else entry.attachment = bytes.slice(); + } + + autoResponse(pair: HibernationAutoResponse | null): void { + this.#autoResponsePair = cloneAutoResponse(pair); + } + + closed(socket: RawWebSocket): void { + this.#entries.delete(socket); + } + + snapshot(): RehydratedWebSocket[] { + return [...this.#entries.values()].map(cloneEntry); + } +} + +function cloneEntry(value: RehydratedWebSocket): MirroredWebSocket { + const entry: MirroredWebSocket = { + socket: value.socket, + tags: [...(value.tags ?? [])], + }; + if (value.attachment !== undefined) entry.attachment = value.attachment.slice(); + if (value.autoResponseTimestamp !== undefined) { + entry.autoResponseTimestamp = value.autoResponseTimestamp; + } + return entry; +} + +function cloneAutoResponse(pair: HibernationAutoResponse | null): HibernationAutoResponse | null { + return pair === null ? null : { ...pair }; +} diff --git a/src/tsconfig.json b/src/tsconfig.json index b46ac60..61cb393 100644 --- a/src/tsconfig.json +++ b/src/tsconfig.json @@ -1,7 +1,7 @@ { "extends": "../tsconfig.base.json", "compilerOptions": { "outDir": "../.tsbuild/index" }, - "files": ["index.ts", "gate.ts", "vite.ts"], + "files": ["browser.ts", "index.ts", "gate.ts", "vite.ts"], "references": [ { "path": "./util" }, { "path": "./io" },