Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/calm-sockets-share.md
Original file line number Diff line number Diff line change
@@ -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.
6 changes: 6 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
6 changes: 4 additions & 2 deletions conformance/browser/actor.worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@
import {
actorScopeBindings,
createActorContainer,
HibernationMirror,
installActorScope,
newRpcSession,
type ActorContainer,
Expand Down Expand Up @@ -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";
Expand Down
6 changes: 4 additions & 2 deletions conformance/node/host.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ import { createNodeSqlProvider } from "../../backends/node-sqlite";
import {
AlarmScheduler,
createActorContainer,
HibernationMirror,
installWebSocketGlobals,
type ActorContainer,
type ActorEntry,
Expand Down Expand Up @@ -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";

// =======================================================================================
Expand Down
42 changes: 2 additions & 40 deletions conformance/websocket-upgrade.ts
Original file line number Diff line number Diff line change
@@ -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);
}
4 changes: 2 additions & 2 deletions examples/extension/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`. |

Expand Down
8 changes: 4 additions & 4 deletions examples/extension/src/worker/actor.worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import {
createActorContainer,
createDurableObjectNamespace,
gateRequestBody,
HibernationMirror,
installActorScope,
newRpcSession,
type ActorContainer,
Expand All @@ -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,
Expand Down Expand Up @@ -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
Expand Down
38 changes: 0 additions & 38 deletions examples/platform-shims/hibernation-mirror.ts

This file was deleted.

44 changes: 0 additions & 44 deletions examples/platform-shims/memory-websocket-pair.ts

This file was deleted.

4 changes: 4 additions & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
1 change: 1 addition & 0 deletions scripts/build-package.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
40 changes: 40 additions & 0 deletions src/browser.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
73 changes: 73 additions & 0 deletions src/browser.ts
Original file line number Diff line number Diff line change
@@ -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<ResponseInit, "webSocket"> & {
webSocket?: UpgradeWebSocket;
};
type UpgradeResponse = Response & { readonly webSocket?: UpgradeWebSocket };
type CloneableRequest = { readonly headers: Headers; clone(): CloneableRequest };

const upgradeRequests = new WeakSet<CloneableRequest>();
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<T extends CloneableRequest>(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"
);
}
2 changes: 2 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
Loading