From 7d2198a9ec406bb407cd4846af341c5963c13ab0 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Mon, 16 Mar 2026 22:58:34 +0300 Subject: [PATCH 01/12] feat/multi-durable-object-support --- docs/2.adapters/cloudflare.md | 20 ++++++++++++++++++ src/adapters/cloudflare.ts | 39 +++++++++++++++++++++++++++++------ src/hooks.ts | 9 ++++++++ 3 files changed, 62 insertions(+), 6 deletions(-) diff --git a/docs/2.adapters/cloudflare.md b/docs/2.adapters/cloudflare.md index f1ea344..5b7cf4a 100644 --- a/docs/2.adapters/cloudflare.md +++ b/docs/2.adapters/cloudflare.md @@ -79,6 +79,25 @@ new_classes = ["$DurableObject"] See [`test/fixture/cloudflare-durable.ts`](https://github.com/h3js/crossws/blob/main/test/fixture/cloudflare-durable.ts) for demo and [`src/adapters/cloudflare.ts`](https://github.com/h3js/crossws/blob/main/src/adapters/cloudflare.ts) for implementation. :: +### Durable Objects +By default, the cloudflare adapter uses a single Durable Object (DO) to handle ALL requests. This behavior will create a bottleneck since DOs are design to scale horizontally and only handle about 1000 requests/s see [DO message throughput limits](https://developers.cloudflare.com/durable-objects/best-practices/rules-of-durable-objects/#message-throughput-limits). To address this, you can use the `useNamespaceAsId` option along with the [`upgrade() hook`](/guide/hooks) to control which DO handle which request. + +There two scenarios for this: + +1. You only enable `useNamespaceAsId`. In this case, a DO will be created to handle the request based on the `URL().pathname`. + +2. You only enable `useNamespaceAsId` and use the upgrade hook in your route to return an object with a `namespace` property. Here you gain full control over the DO creation. The namespace to get the DO instance id. + +> [!NOTE] +> When you enable the `useNamespaceAsId` option, your `upgrade()` hook will run twice!. First it will run on the worker to check if you have returned a `namespace`. Then it will run on the DO once the connection is passed to it. You can use the second argument of the `upgrade()` hook which contains the upgrade context. +>```ts +>type UpgradeContext = { +> cf?: { +> runtime: "worker" | "DO"; +> }; +>}; +>``` + ### Adapter options > [!NOTE] @@ -87,4 +106,5 @@ See [`test/fixture/cloudflare-durable.ts`](https://github.com/h3js/crossws/blob/ - `bindingName`: Durable Object binding name from environment (default: `$DurableObject`). - `instanceName`: Durable Object instance name (default: `crossws`). +- `useNamespaceAsId`: When this is set to `true`, each peer namespace will have their its own Durable Object (default: false). - `resolveDurableStub`: Custom function that resolves Durable Object binding to handle the WebSocket upgrade. This option will override `bindingName` and `instanceName`. diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index 635856d..8aebf3d 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -5,7 +5,7 @@ import type * as web from "../../types/web.ts"; import { env as cfGlobalEnv } from "cloudflare:workers"; import { toBufferLike } from "../utils.ts"; import { adapterUtils, getPeers } from "../adapter.ts"; -import { AdapterHookable } from "../hooks.ts"; +import { AdapterHookable, type UpgradeContext } from "../hooks.ts"; import { Message } from "../message.ts"; import { Peer, type PeerContext } from "../peer.ts"; import { StubRequest } from "../_request.ts"; @@ -40,10 +40,21 @@ export interface CloudflareOptions extends AdapterOptions { */ instanceName?: string; + /** + * Create durable object for each namespace. + * + * **Note:** This option will be ignored if `resolveDurableStub` is provided. + * + * **Note:** This option will cause the upgrade hook to run twice!. + * + * @default false + */ + useNamespaceAsId?: boolean; + /** * Custom function that resolves Durable Object binding to handle the WebSocket upgrade. * - * **Note:** This option will override `bindingName` and `instanceName`. + * **Note:** This option will override `bindingName`, `instanceName` and `useNamespaceAsStubId`. */ resolveDurableStub?: ResolveDurableStub; } @@ -62,14 +73,27 @@ const cloudflareAdapter: Adapter< const resolveDurableStub: ResolveDurableStub = opts.resolveDurableStub || - ((_req, env: any, _context): WSDurableObjectStub | undefined => { + (async ( + _req, + env: any, + _context, + ): Promise => { const bindingName = opts.bindingName || "$DurableObject"; const binding = (env || cfGlobalEnv)[ bindingName ] as CF.DurableObjectNamespace; if (binding) { - const instanceId = binding.idFromName(opts.instanceName || "crossws"); - return binding.get(instanceId); + let instanceName = opts.instanceName || "crossws"; + + if (opts.useNamespaceAsId) { + const { namespace } = await hooks.upgrade( + _req as unknown as Request, + { cf: { runtime: "worker" } }, + ); + if (namespace) instanceName = namespace; + } + + return binding.get(binding.idFromName(instanceName)); } }); @@ -92,7 +116,9 @@ const cloudflareAdapter: Adapter< // [Fallback] Upgrade request in same Worker const { upgradeHeaders, endResponse, context, namespace } = - await hooks.upgrade(request as unknown as Request); + await hooks.upgrade(request as unknown as Request, { + cf: { runtime: "worker" }, + }); if (endResponse) { return endResponse as unknown as Response; } @@ -148,6 +174,7 @@ const cloudflareAdapter: Adapter< handleDurableUpgrade: async (obj, request) => { const { upgradeHeaders, endResponse, namespace } = await hooks.upgrade( request as Request, + { cf: { runtime: "DO" } }, ); if (endResponse) { return endResponse; diff --git a/src/hooks.ts b/src/hooks.ts index c665abe..703961d 100644 --- a/src/hooks.ts +++ b/src/hooks.ts @@ -43,6 +43,7 @@ export class AdapterHookable { async upgrade( request: Request & { readonly context?: Record }, + upgradeContext?: UpgradeContext, ): Promise<{ context: PeerContext; namespace: string; @@ -58,6 +59,7 @@ export class AdapterHookable { const res = await this.callHook( "upgrade", request as Request & { context?: PeerContext }, + upgradeContext, ); if (!res) { return { context, namespace }; @@ -112,6 +114,12 @@ export type MaybePromise = T | Promise; export type UpgradeError = Response | { readonly response: Response }; +export type UpgradeContext = { + cf?: { + runtime: "worker" | "DO"; + }; +}; + export interface Hooks { /** * Upgrading a request to a WebSocket connection. @@ -128,6 +136,7 @@ export interface Hooks { request: Request & { readonly context?: Record; }, + context?: UpgradeContext, ) => MaybePromise< | { headers?: HeadersInit; namespace?: string; context?: PeerContext } | Response From c0a0adbb450a10afab61484461558dc8d1cb6121 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Mon, 16 Mar 2026 23:43:25 +0300 Subject: [PATCH 02/12] fixed typos --- docs/2.adapters/cloudflare.md | 8 ++++---- src/adapters/cloudflare.ts | 2 +- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/docs/2.adapters/cloudflare.md b/docs/2.adapters/cloudflare.md index 5b7cf4a..64fa8cf 100644 --- a/docs/2.adapters/cloudflare.md +++ b/docs/2.adapters/cloudflare.md @@ -79,10 +79,10 @@ new_classes = ["$DurableObject"] See [`test/fixture/cloudflare-durable.ts`](https://github.com/h3js/crossws/blob/main/test/fixture/cloudflare-durable.ts) for demo and [`src/adapters/cloudflare.ts`](https://github.com/h3js/crossws/blob/main/src/adapters/cloudflare.ts) for implementation. :: -### Durable Objects +## Durable Objects By default, the cloudflare adapter uses a single Durable Object (DO) to handle ALL requests. This behavior will create a bottleneck since DOs are design to scale horizontally and only handle about 1000 requests/s see [DO message throughput limits](https://developers.cloudflare.com/durable-objects/best-practices/rules-of-durable-objects/#message-throughput-limits). To address this, you can use the `useNamespaceAsId` option along with the [`upgrade() hook`](/guide/hooks) to control which DO handle which request. -There two scenarios for this: +There are two scenarios for this: 1. You only enable `useNamespaceAsId`. In this case, a DO will be created to handle the request based on the `URL().pathname`. @@ -98,7 +98,7 @@ There two scenarios for this: >}; >``` -### Adapter options +## Adapter options > [!NOTE] > By default, crossws uses the durable object class `$DurableObject` from `env` with an instance named `crossws`. @@ -106,5 +106,5 @@ There two scenarios for this: - `bindingName`: Durable Object binding name from environment (default: `$DurableObject`). - `instanceName`: Durable Object instance name (default: `crossws`). -- `useNamespaceAsId`: When this is set to `true`, each peer namespace will have their its own Durable Object (default: false). +- `useNamespaceAsId`: When set to `true`, each peer namespace gets its own Durable Object (default: `false`). - `resolveDurableStub`: Custom function that resolves Durable Object binding to handle the WebSocket upgrade. This option will override `bindingName` and `instanceName`. diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index 8aebf3d..d6698f8 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -54,7 +54,7 @@ export interface CloudflareOptions extends AdapterOptions { /** * Custom function that resolves Durable Object binding to handle the WebSocket upgrade. * - * **Note:** This option will override `bindingName`, `instanceName` and `useNamespaceAsStubId`. + * **Note:** This option will override `bindingName`, `instanceName` and `useNamespaceAsId`. */ resolveDurableStub?: ResolveDurableStub; } From 4b9e14d8fd635a4729ee93bc4d1d83d29fb4b57b Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Tue, 17 Mar 2026 04:49:51 +0300 Subject: [PATCH 03/12] handling undefined request --- src/adapters/cloudflare.ts | 66 +++++++++++++++++++++++--------------- 1 file changed, 41 insertions(+), 25 deletions(-) diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index d6698f8..1303090 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -19,6 +19,7 @@ type ResolveDurableStub = ( req: CF.Request | undefined, env: unknown, context: CF.ExecutionContext | undefined, + namespace?: string, ) => WSDurableObjectStub | undefined | Promise; export interface CloudflareOptions extends AdapterOptions { @@ -71,31 +72,40 @@ const cloudflareAdapter: Adapter< Set >(); - const resolveDurableStub: ResolveDurableStub = - opts.resolveDurableStub || - (async ( - _req, - env: any, - _context, - ): Promise => { - const bindingName = opts.bindingName || "$DurableObject"; - const binding = (env || cfGlobalEnv)[ - bindingName - ] as CF.DurableObjectNamespace; - if (binding) { - let instanceName = opts.instanceName || "crossws"; - - if (opts.useNamespaceAsId) { - const { namespace } = await hooks.upgrade( - _req as unknown as Request, - { cf: { runtime: "worker" } }, - ); - if (namespace) instanceName = namespace; - } - - return binding.get(binding.idFromName(instanceName)); + const defaultDurableStubResolver: ResolveDurableStub = async ( + req, + env: any, + _context, + explicitNamespace, + ) => { + const bindingName = opts.bindingName || "$DurableObject"; + const binding = (env || cfGlobalEnv)[ + bindingName + ] as CF.DurableObjectNamespace; + + if (!binding) { + return undefined; + } + + // Determine the ID name logic: + // 1. Use explicitNamespace if provided (e.g. from publish(..., { namespace })) + // 2. If useNamespaceAsId is true and we have a request, run the upgrade hook + // 3. Fallback to instanceName + let instanceName = explicitNamespace || opts.instanceName || "crossws"; + + if (!explicitNamespace && opts.useNamespaceAsId && req) { + const { namespace } = await hooks.upgrade(req as unknown as Request, { + cf: { runtime: "worker" }, + }); + if (namespace) { + instanceName = namespace; } - }); + } + + return binding.get(binding.idFromName(instanceName)); + }; + const resolveDurableStub: ResolveDurableStub = + opts.resolveDurableStub || defaultDurableStubResolver; const { publish: durablePublish, ...utils } = adapterUtils(globalPeers); @@ -216,7 +226,13 @@ const cloudflareAdapter: Adapter< return durablePublish(topic, data, opts); }, publish: async (topic, data, opts) => { - const stub = await resolveDurableStub(undefined, cfGlobalEnv, undefined); + const stub = await resolveDurableStub( + undefined, + cfGlobalEnv, + undefined, + opts?.namespace, + ); + if (!stub) { throw new Error("[crossws] Durable Object binding cannot be resolved."); } From 2629572e7350eae7b2fa6436f5ff32a055e0add1 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 05:27:12 +0300 Subject: [PATCH 04/12] fixed cloudflare I/O error --- src/adapters/cloudflare.ts | 65 +++++++++++++++++++++++--------------- 1 file changed, 40 insertions(+), 25 deletions(-) diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index 1303090..b6364ab 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -4,8 +4,7 @@ import type { AdapterOptions, AdapterInstance, Adapter } from "../adapter.ts"; import type * as web from "../../types/web.ts"; import { env as cfGlobalEnv } from "cloudflare:workers"; import { toBufferLike } from "../utils.ts"; -import { adapterUtils, getPeers } from "../adapter.ts"; -import { AdapterHookable, type UpgradeContext } from "../hooks.ts"; +import { AdapterHookable } from "../hooks.ts"; import { Message } from "../message.ts"; import { Peer, type PeerContext } from "../peer.ts"; import { StubRequest } from "../_request.ts"; @@ -104,13 +103,14 @@ const cloudflareAdapter: Adapter< return binding.get(binding.idFromName(instanceName)); }; + const resolveDurableStub: ResolveDurableStub = opts.resolveDurableStub || defaultDurableStubResolver; - const { publish: durablePublish, ...utils } = adapterUtils(globalPeers); - return { - ...utils, + // Returns an empty Map(). Accessing this object across different requests or Durable Objects on Cloudflare triggers I/O errors, + // rendering it non-functional in those contexts. Maintained solely for backward compatibility. + peers: new Map(), handleUpgrade: async (request, cfEnv, cfCtx) => { // Upgrade request with Durable Object binding const stub = await resolveDurableStub( @@ -132,10 +132,7 @@ const cloudflareAdapter: Adapter< if (endResponse) { return endResponse as unknown as Response; } - const peers = getPeers( - globalPeers, - namespace, - ) as Set; + const pair = new WebSocketPair() as unknown as [ CF.WebSocket, CF.WebSocket, @@ -144,7 +141,6 @@ const cloudflareAdapter: Adapter< const server = pair[1]; const peer = new CloudflareFallbackPeer({ ws: client, - peers, wsServer: server, request: request as unknown as Request, cfEnv, @@ -152,7 +148,6 @@ const cloudflareAdapter: Adapter< context, namespace, }); - peers.add(peer); server.accept(); hooks.callHook("open", peer); server.addEventListener("message", (event) => { @@ -163,11 +158,9 @@ const cloudflareAdapter: Adapter< ); }); server.addEventListener("error", (event) => { - peers.delete(peer); hooks.callHook("error", peer, new WSError(event.error)); }); server.addEventListener("close", (event) => { - peers.delete(peer); hooks.callHook("close", peer, event); server.close(); }); @@ -190,8 +183,6 @@ const cloudflareAdapter: Adapter< return endResponse; } - const peers = getPeers(globalPeers, namespace); - const pair = new WebSocketPair(); const client = pair[0]; const server = pair[1]; @@ -201,7 +192,7 @@ const cloudflareAdapter: Adapter< request, namespace, ); - peers.add(peer); + (obj as DurableObjectPub).ctx.acceptWebSocket(server); await hooks.callHook("open", peer); @@ -217,13 +208,14 @@ const cloudflareAdapter: Adapter< }, handleDurableClose: async (obj, ws, code, reason, wasClean) => { const peer = CloudflareDurablePeer._restore(obj, ws as CF.WebSocket); - const peers = getPeers(globalPeers, peer.namespace); - peers.delete(peer); const details = { code, reason, wasClean }; await hooks.callHook("close", peer, details); }, handleDurablePublish: async (_obj, topic, data, opts) => { - return durablePublish(topic, data, opts); + const peers = getDurablePeers(_obj as DurableObjectPub, topic); + for (const peer of peers) { + peer.send(data); + } }, publish: async (topic, data, opts) => { const stub = await resolveDurableStub( @@ -252,6 +244,25 @@ export default cloudflareAdapter; // --- peer --- +function getDurablePeers( + obj: DurableObjectPub, + topic?: string, +): CloudflareDurablePeer[] { + const peers: CloudflareDurablePeer[] = []; + + const websockets = obj.ctx.getWebSockets() as unknown as AugmentedWebSocket[]; + for (const ws of websockets) { + const state = getAttachedState(ws); + if (topic && state.t && !state.t.has(topic)) { + continue; + } + + const peer = CloudflareDurablePeer._restore(obj, ws); + peers.push(peer); + } + return peers; +} + class CloudflareDurablePeer extends Peer<{ ws: AugmentedWebSocket; request: Request; @@ -292,10 +303,10 @@ class CloudflareDurablePeer extends Peer<{ } const dataBuff = toBufferLike(data); for (const ws of websockets) { - if (ws === this._internal.ws) { + const state = getAttachedState(ws); + if (state.i === this.id) { continue; } - const state = getAttachedState(ws); if (state.t?.has(topic)) { ws.send(dataBuff); } @@ -317,11 +328,13 @@ class CloudflareDurablePeer extends Peer<{ return peer; } const state = (ws.deserializeAttachment() || {}) as AttachedState; + const peerNamespace = + namespace || state.n || ""; /* later throws error if empty */ peer = ws._crosswsPeer = new CloudflareDurablePeer({ ws: ws as CF.WebSocket, request: (request as Request | undefined) || new StubRequest(state.u || ""), - namespace: namespace || state.n || "" /* later throws error if empty */, + namespace: peerNamespace, durable: durable as DurableObjectPub, }); if (state.i) { @@ -331,6 +344,9 @@ class CloudflareDurablePeer extends Peer<{ state.u = request.url; } state.i = peer.id; + + state.n = peerNamespace; + setAttachedState(ws, state); return peer; } @@ -339,7 +355,6 @@ class CloudflareDurablePeer extends Peer<{ class CloudflareFallbackPeer extends Peer<{ ws: CF.WebSocket; request: Request; - peers: Set; wsServer: CF.WebSocket; cfEnv: unknown; cfCtx: CF.ExecutionContext; @@ -399,8 +414,8 @@ type AttachedState = { i?: string; /** Request url */ u?: string; - /** Connection namespace */ - n?: string; + /** Connection namespace mandatory! */ + n: string; }; export interface CloudflareDurableAdapter extends AdapterInstance { From 0de349aa8c63cb0690c6d4bfc0a2c0ea8f5bb258 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 05:46:07 +0300 Subject: [PATCH 05/12] fixed namespaces ignored in handleDurablePublish() --- src/adapters/cloudflare.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index b6364ab..a7b43e2 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -214,6 +214,10 @@ const cloudflareAdapter: Adapter< handleDurablePublish: async (_obj, topic, data, opts) => { const peers = getDurablePeers(_obj as DurableObjectPub, topic); for (const peer of peers) { + // single Durable Object with multiple namesapces + if (peer.namespace !== opts.namespace) { + continue; + } peer.send(data); } }, From 6691288e3e40b78b8f79d92cb1f3754e9c7fe857 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 05:46:27 +0300 Subject: [PATCH 06/12] fixed typo --- src/adapters/cloudflare.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index a7b43e2..6d103f3 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -214,7 +214,7 @@ const cloudflareAdapter: Adapter< handleDurablePublish: async (_obj, topic, data, opts) => { const peers = getDurablePeers(_obj as DurableObjectPub, topic); for (const peer of peers) { - // single Durable Object with multiple namesapces + // single Durable Object with multiple namespaces if (peer.namespace !== opts.namespace) { continue; } From f2245b7372b5be8f565de9dfe3874ba9efa70ac6 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 05:47:39 +0300 Subject: [PATCH 07/12] added type --- src/adapters/cloudflare.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index 6d103f3..3c134ec 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -450,7 +450,7 @@ export interface CloudflareDurableAdapter extends AdapterInstance { obj: DurableObject, topic: string, data: unknown, - opts: any, + opts: Record & { namespace?: string }, ) => Promise; handleDurableClose( From 8f434ef0e91588a87dd86663bea45e92e0fae44d Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 05:51:03 +0300 Subject: [PATCH 08/12] exposed getDurablePeers() --- src/adapters/cloudflare.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index 3c134ec..b13ecc0 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -111,6 +111,7 @@ const cloudflareAdapter: Adapter< // Returns an empty Map(). Accessing this object across different requests or Durable Objects on Cloudflare triggers I/O errors, // rendering it non-functional in those contexts. Maintained solely for backward compatibility. peers: new Map(), + getDurablePeers, handleUpgrade: async (request, cfEnv, cfCtx) => { // Upgrade request with Durable Object binding const stub = await resolveDurableStub( From d9de0f06498b3bf26b2d44a55fe10e2094b57147 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 05:57:29 +0300 Subject: [PATCH 09/12] updated cloudflare docs --- docs/2.adapters/cloudflare.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/docs/2.adapters/cloudflare.md b/docs/2.adapters/cloudflare.md index 64fa8cf..c34fa1b 100644 --- a/docs/2.adapters/cloudflare.md +++ b/docs/2.adapters/cloudflare.md @@ -98,6 +98,9 @@ There are two scenarios for this: >}; >``` +> [!WARNING] +> When using the Cloudflare adapter, the `peers` property of the adapter is always set to an empty `Map()`. If you need access to the list of connected peers within a DO use the `getDurablePeers()` function instead. `getDurablePeers()` can only be used inside the `$DurableObject` class since it requires the DO instance. + ## Adapter options > [!NOTE] From 0531f0326dae29613b93bec255ac5b4c96a703cf Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 06:05:44 +0300 Subject: [PATCH 10/12] fixed bugs --- src/adapters/cloudflare.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index b13ecc0..4ce0330 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -216,7 +216,7 @@ const cloudflareAdapter: Adapter< const peers = getDurablePeers(_obj as DurableObjectPub, topic); for (const peer of peers) { // single Durable Object with multiple namespaces - if (peer.namespace !== opts.namespace) { + if (opts && peer.namespace !== opts.namespace) { continue; } peer.send(data); @@ -258,7 +258,7 @@ function getDurablePeers( const websockets = obj.ctx.getWebSockets() as unknown as AugmentedWebSocket[]; for (const ws of websockets) { const state = getAttachedState(ws); - if (topic && state.t && !state.t.has(topic)) { + if (topic && !state.t?.has(topic)) { continue; } From 0f457716ebd2f9812e82ab95c4a5ffedbb82d381 Mon Sep 17 00:00:00 2001 From: Abdullah Azbah Date: Wed, 18 Mar 2026 06:10:39 +0300 Subject: [PATCH 11/12] small changes to docs --- docs/2.adapters/cloudflare.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/2.adapters/cloudflare.md b/docs/2.adapters/cloudflare.md index c34fa1b..08bfd75 100644 --- a/docs/2.adapters/cloudflare.md +++ b/docs/2.adapters/cloudflare.md @@ -110,4 +110,4 @@ There are two scenarios for this: - `bindingName`: Durable Object binding name from environment (default: `$DurableObject`). - `instanceName`: Durable Object instance name (default: `crossws`). - `useNamespaceAsId`: When set to `true`, each peer namespace gets its own Durable Object (default: `false`). -- `resolveDurableStub`: Custom function that resolves Durable Object binding to handle the WebSocket upgrade. This option will override `bindingName` and `instanceName`. +- `resolveDurableStub`: Custom function that resolves Durable Object binding to handle the WebSocket upgrade. This option will override `bindingName`, `useNamespaceAsId`, and `instanceName`. From 2c5664581c1918c1f678367dbb7355dd5167ea87 Mon Sep 17 00:00:00 2001 From: Pooya Parsa Date: Fri, 3 Jul 2026 09:30:52 +0000 Subject: [PATCH 12/12] wip --- docs/2.adapters/cloudflare.md | 11 ++++++++- src/adapters/cloudflare.ts | 37 +++++++++++++++++++++++++----- test/fixture/cloudflare-durable.ts | 11 +++++++++ 3 files changed, 52 insertions(+), 7 deletions(-) diff --git a/docs/2.adapters/cloudflare.md b/docs/2.adapters/cloudflare.md index 08bfd75..057e826 100644 --- a/docs/2.adapters/cloudflare.md +++ b/docs/2.adapters/cloudflare.md @@ -99,7 +99,16 @@ There are two scenarios for this: >``` > [!WARNING] -> When using the Cloudflare adapter, the `peers` property of the adapter is always set to an empty `Map()`. If you need access to the list of connected peers within a DO use the `getDurablePeers()` function instead. `getDurablePeers()` can only be used inside the `$DurableObject` class since it requires the DO instance. +> When using Durable Objects, the adapter's `peers` property never contains the Durable Object peers (it only tracks the in-Worker fallback path). If you need access to the list of connected peers within a DO use the `getDurablePeers()` function instead. `getDurablePeers()` can only be used inside the `$DurableObject` class since it requires the DO instance. +>```ts +>export class $DurableObject extends DurableObject { +> // Pass `this` since it requires the Durable Object instance. +> // Optionally pass a topic to only list peers subscribed to it. +> listPeers() { +> return ws.getDurablePeers(this).map((peer) => peer.id); +> } +>} +>``` ## Adapter options diff --git a/src/adapters/cloudflare.ts b/src/adapters/cloudflare.ts index 1e46962..5155bba 100644 --- a/src/adapters/cloudflare.ts +++ b/src/adapters/cloudflare.ts @@ -4,6 +4,7 @@ import type { AdapterOptions, AdapterInstance, Adapter } from "../adapter.ts"; import type * as web from "../../types/web.ts"; import { env as cfGlobalEnv } from "cloudflare:workers"; import { toBufferLike } from "../utils.ts"; +import { getPeers } from "../adapter.ts"; import { AdapterHookable } from "../hooks.ts"; import { Message } from "../message.ts"; import { Peer, type PeerContext } from "../peer.ts"; @@ -63,7 +64,12 @@ export interface CloudflareOptions extends AdapterOptions { const cloudflareAdapter: Adapter = (opts = {}) => { const hooks = new AdapterHookable(opts); - const globalPeers = new Map>(); + + // Tracks peers for the in-Worker fallback path only (single isolate, no + // Durable Object). Durable Object peers are intentionally NOT tracked here: + // holding peer references across requests/DOs on Cloudflare triggers I/O + // errors. Use `getDurablePeers()` from within the Durable Object instead. + const globalPeers = new Map>(); const defaultDurableStubResolver: ResolveDurableStub = async ( req, @@ -100,9 +106,9 @@ const cloudflareAdapter: Adapter = opts.resolveDurableStub || defaultDurableStubResolver; return { - // Returns an empty Map(). Accessing this object across different requests or Durable Objects on Cloudflare triggers I/O errors, - // rendering it non-functional in those contexts. Maintained solely for backward compatibility. - peers: new Map(), + // Only the in-Worker fallback peers are exposed here. Durable Object peers + // are not (see `globalPeers` above); use `getDurablePeers()` for those. + peers: globalPeers, getDurablePeers, handleUpgrade: async (request, cfEnv, cfCtx) => { // Upgrade request with Durable Object binding @@ -120,11 +126,13 @@ const cloudflareAdapter: Adapter = return endResponse as unknown as Response; } + const peers = getPeers(globalPeers, namespace); const pair = new WebSocketPair() as unknown as [CF.WebSocket, CF.WebSocket]; const client = pair[0]; const server = pair[1]; const peer = new CloudflareFallbackPeer({ ws: client, + peers, wsServer: server, request: request as unknown as Request, cfEnv, @@ -132,6 +140,7 @@ const cloudflareAdapter: Adapter = context, namespace, }); + peers.add(peer); server.accept(); hooks.callHook("open", peer); server.addEventListener("message", (event) => { @@ -142,9 +151,11 @@ const cloudflareAdapter: Adapter = ); }); server.addEventListener("error", (event) => { + peers.delete(peer); hooks.callHook("error", peer, new WSError(event.error)); }); server.addEventListener("close", (event) => { + peers.delete(peer); hooks.callHook("close", peer, event); server.close(); }); @@ -197,8 +208,10 @@ const cloudflareAdapter: Adapter = handleDurablePublish: async (_obj, topic, data, opts) => { const peers = getDurablePeers(_obj as DurableObjectPub, topic); for (const peer of peers) { - // single Durable Object with multiple namespaces - if (opts && peer.namespace !== opts.namespace) { + // When a namespace is given, scope the publish to a single namespace + // (single Durable Object hosting multiple namespaces). Without it, + // publish to every namespace subscribed to the topic. + if (opts?.namespace && peer.namespace !== opts.namespace) { continue; } peer.send(data); @@ -328,6 +341,7 @@ class CloudflareDurablePeer extends Peer<{ class CloudflareFallbackPeer extends Peer<{ ws: CF.WebSocket; request: Request; + peers: Set; wsServer: CF.WebSocket; cfEnv: unknown; cfCtx: CF.ExecutionContext; @@ -390,6 +404,17 @@ type AttachedState = { }; export interface CloudflareDurableAdapter extends AdapterInstance { + /** + * List the peers connected to a Durable Object instance, optionally filtered + * by `topic`. + * + * **Note:** Must be called from within the `$DurableObject` class (e.g. + * `ws.getDurablePeers(this)`) since it relies on the Durable Object context. + * The adapter-level `peers` map only tracks the in-Worker fallback path and + * never contains Durable Object peers. + */ + getDurablePeers(obj: DurableObject, topic?: string): Peer[]; + handleUpgrade( req: Request | CF.Request, env: unknown, diff --git a/test/fixture/cloudflare-durable.ts b/test/fixture/cloudflare-durable.ts index 1e5b197..5989626 100644 --- a/test/fixture/cloudflare-durable.ts +++ b/test/fixture/cloudflare-durable.ts @@ -11,6 +11,13 @@ export default { env: Record, context: ExecutionContext, ): Promise { + // The adapter-level `peers` map is always empty on Cloudflare; peers live + // inside the Durable Object and must be enumerated from within it. + if (new URL(request.url).pathname === "/peers") { + const stub = env.$DurableObject.get(env.$DurableObject.idFromName("crossws")); + return Response.json({ peers: await stub.webSocketPeers() }); + } + const response = handleDemoRoutes(ws, request); if (response) { return response; @@ -40,6 +47,10 @@ export class $DurableObject extends DurableObject { return ws.handleDurablePublish(this, topic, message, opts); } + webSocketPeers() { + return ws.getDurablePeers(this).map((peer) => `${peer.namespace}:${peer.id}`); + } + override async webSocketMessage(client: WebSocket, message: ArrayBuffer | string): Promise { return ws.handleDurableMessage(this, client, message); }